peer.Peer is now an interface

Juan Batiz-Benet committed
Oct 20, 2014 at 03:26 UTC
9ca87fbb936729881d87ba0c36beb03b9a00ade8
65 files changed
+538
-459
core/commands/publish.go
+1
-1
@@ -35,7 +35,7 @@ func Publish(n *core.IpfsNode, args []string, opts map[string]interface{}, out i
35
}
36
37
// later, n.Keychain.Get(name).PrivKey
38
- k := n.Identity.PrivKey
38
+ k := n.Identity.PrivKey()
39
40
pub := nsys.NewRoutingPublisher(n.Routing)
41
err := pub.Publish(k, ref)
core/commands/resolve.go
+1
-1
@@ -19,7 +19,7 @@ func Resolve(n *core.IpfsNode, args []string, opts map[string]interface{}, out i
19
if n.Identity == nil {
20
return errors.New("Identity not loaded!")
21
}
22
- name = n.Identity.ID.String()
22
+ name = n.Identity.ID().String()
23
24
default:
25
return fmt.Errorf("Publish expects 1 or 2 args; got %d.", len(args))
core/core.go
+2
-2
@@ -36,7 +36,7 @@ type IpfsNode struct {
36
Config *config.Config
37
38
// the local node's identity
39
- Identity *peer.Peer
39
+ Identity peer.Peer
40
41
// storage for other Peer instances
42
Peerstore peer.Peerstore
@@ -177,7 +177,7 @@ func NewIpfsNode(cfg *config.Config, online bool) (*IpfsNode, error) {
177
}, nil
178
}
179
180
-func initIdentity(cfg *config.Config, peers peer.Peerstore, online bool) (*peer.Peer, error) {
180
+func initIdentity(cfg *config.Config, peers peer.Peerstore, online bool) (peer.Peer, error) {
181
if cfg.Identity.PeerID == "" {
182
return nil, errors.New("Identity was not set in config (was ipfs init run?)")
183
}
core/mock.go
+6
-8
@@ -8,7 +8,7 @@ import (
8
mdag "github.com/jbenet/go-ipfs/merkledag"
9
nsys "github.com/jbenet/go-ipfs/namesys"
10
path "github.com/jbenet/go-ipfs/path"
11
- "github.com/jbenet/go-ipfs/peer"
11
+ peer "github.com/jbenet/go-ipfs/peer"
12
mdht "github.com/jbenet/go-ipfs/routing/mock"
13
)
14
@@ -16,20 +16,18 @@ import (
16
func NewMockNode() (*IpfsNode, error) {
17
nd := new(IpfsNode)
18
19
- //Generate Identity
20
- nd.Peerstore = peer.NewPeerstore()
21
- var err error
22
- nd.Identity, err = nd.Peerstore.Get(peer.ID("TESTING"))
19
+ // Generate Identity
20
+ sk, pk, err := ci.GenerateKeyPair(ci.RSA, 1024)
21
if err != nil {
22
return nil, err
23
}
24
27
- pk, sk, err := ci.GenerateKeyPair(ci.RSA, 1024)
25
+ nd.Identity, err = peer.WithKeyPair(sk, pk)
26
if err != nil {
27
return nil, err
28
}
31
- nd.Identity.PrivKey = pk
32
- nd.Identity.PubKey = sk
29
+ nd.Peerstore = peer.NewPeerstore()
30
+ nd.Peerstore.Put(nd.Identity)
31
32
// Temp Datastore
33
dstore := ds.NewMapDatastore()
crypto/spipe/handshake.go
+4
-4
@@ -53,7 +53,7 @@ func (s *SecurePipe) handshake() error {
53
}
54
55
log.Debug("handshake: %s <--> %s", s.local, s.remote)
56
- myPubKey, err := s.local.PubKey.Bytes()
56
+ myPubKey, err := s.local.PubKey().Bytes()
57
if err != nil {
58
return err
59
}
@@ -132,7 +132,7 @@ func (s *SecurePipe) handshake() error {
132
exPacket := new(Exchange)
133
134
exPacket.Epubkey = epubkey
135
- exPacket.Signature, err = s.local.PrivKey.Sign(handshake.Bytes())
135
+ exPacket.Signature, err = s.local.PrivKey().Sign(handshake.Bytes())
136
if err != nil {
137
return err
138
}
@@ -167,7 +167,7 @@ func (s *SecurePipe) handshake() error {
167
theirHandshake.Write(exchangeResp.GetEpubkey())
168
169
// u.POut("Remote Peer Identified as %s\n", s.remote)
170
- ok, err := s.remote.PubKey.Verify(theirHandshake.Bytes(), exchangeResp.GetSignature())
170
+ ok, err := s.remote.PubKey().Verify(theirHandshake.Bytes(), exchangeResp.GetSignature())
171
if err != nil {
172
return err
173
}
@@ -340,7 +340,7 @@ func selectBest(myPrefs, theirPrefs string) (string, error) {
340
// getOrConstructPeer attempts to fetch a peer from a peerstore.
341
// if succeeds, verify ID and PubKey match.
342
// else, construct it.
343
-func getOrConstructPeer(peers peer.Peerstore, rpk ci.PubKey) (*peer.Peer, error) {
343
+func getOrConstructPeer(peers peer.Peerstore, rpk ci.PubKey) (peer.Peer, error) {
344
345
rid, err := peer.IDFromPubKey(rpk)
346
if err != nil {
crypto/spipe/pipe.go
+5
-5
@@ -18,8 +18,8 @@ type SecurePipe struct {
18
Duplex
19
insecure Duplex
20
21
- local *peer.Peer
22
- remote *peer.Peer
21
+ local peer.Peer
22
+ remote peer.Peer
23
peers peer.Peerstore
24
25
params params
@@ -33,7 +33,7 @@ type params struct {
33
}
34
35
// NewSecurePipe constructs a pipe with channels of a given buffer size.
36
-func NewSecurePipe(ctx context.Context, bufsize int, local *peer.Peer,
36
+func NewSecurePipe(ctx context.Context, bufsize int, local peer.Peer,
37
peers peer.Peerstore, insecure Duplex) (*SecurePipe, error) {
38
39
ctx, cancel := context.WithCancel(ctx)
@@ -60,12 +60,12 @@ func NewSecurePipe(ctx context.Context, bufsize int, local *peer.Peer,
60
}
61
62
// LocalPeer retrieves the local peer.
63
-func (s *SecurePipe) LocalPeer() *peer.Peer {
63
+func (s *SecurePipe) LocalPeer() peer.Peer {
64
return s.local
65
}
66
67
// RemotePeer retrieves the local peer.
68
-func (s *SecurePipe) RemotePeer() *peer.Peer {
68
+func (s *SecurePipe) RemotePeer() peer.Peer {
69
return s.remote
70
}
71
diagnostics/diag.go
+12
-11
@@ -26,14 +26,14 @@ const ResponseTimeout = time.Second * 10
26
type Diagnostics struct {
27
network net.Network
28
sender net.Sender
29
- self *peer.Peer
29
+ self peer.Peer
30
31
diagLock sync.Mutex
32
diagMap map[string]time.Time
33
birth time.Time
34
}
35
36
-func NewDiagnostics(self *peer.Peer, inet net.Network, sender net.Sender) *Diagnostics {
36
+func NewDiagnostics(self peer.Peer, inet net.Network, sender net.Sender) *Diagnostics {
37
return &Diagnostics{
38
network: inet,
39
sender: sender,
@@ -67,20 +67,21 @@ func (di *DiagInfo) Marshal() []byte {
67
return b
68
}
69
70
-func (d *Diagnostics) getPeers() []*peer.Peer {
70
+func (d *Diagnostics) getPeers() []peer.Peer {
71
return d.network.GetPeerList()
72
}
73
74
func (d *Diagnostics) getDiagInfo() *DiagInfo {
75
di := new(DiagInfo)
76
di.CodeVersion = "github.com/jbenet/go-ipfs"
77
- di.ID = d.self.ID.Pretty()
77
+ di.ID = d.self.ID().Pretty()
78
di.LifeSpan = time.Since(d.birth)
79
di.Keys = nil // Currently no way to query datastore
80
di.BwIn, di.BwOut = d.network.GetBandwidthTotals()
81
82
for _, p := range d.getPeers() {
83
- di.Connections = append(di.Connections, connDiagInfo{p.GetLatency(), p.ID.Pretty()})
83
+ d := connDiagInfo{p.GetLatency(), p.ID().Pretty()}
84
+ di.Connections = append(di.Connections, d)
85
}
86
return di
87
}
@@ -116,7 +117,7 @@ func (d *Diagnostics) GetDiagnostic(timeout time.Duration) ([]*DiagInfo, error)
117
for _, p := range peers {
118
log.Debug("Sending getDiagnostic to: %s", p)
119
sends++
119
- go func(p *peer.Peer) {
120
+ go func(p peer.Peer) {
121
data, err := d.getDiagnosticFromPeer(ctx, p, pmes)
122
if err != nil {
123
log.Error("GetDiagnostic error: %v", err)
@@ -155,7 +156,7 @@ func AppendDiagnostics(data []byte, cur []*DiagInfo) []*DiagInfo {
156
}
157
158
// TODO: this method no longer needed.
158
-func (d *Diagnostics) getDiagnosticFromPeer(ctx context.Context, p *peer.Peer, mes *Message) ([]byte, error) {
159
+func (d *Diagnostics) getDiagnosticFromPeer(ctx context.Context, p peer.Peer, mes *Message) ([]byte, error) {
160
rpmes, err := d.sendRequest(ctx, p, mes)
161
if err != nil {
162
return nil, err
@@ -169,7 +170,7 @@ func newMessage(diagID string) *Message {
170
return pmes
171
}
172
172
-func (d *Diagnostics) sendRequest(ctx context.Context, p *peer.Peer, pmes *Message) (*Message, error) {
173
+func (d *Diagnostics) sendRequest(ctx context.Context, p peer.Peer, pmes *Message) (*Message, error) {
174
175
mes, err := msg.FromObject(p, pmes)
176
if err != nil {
@@ -197,7 +198,7 @@ func (d *Diagnostics) sendRequest(ctx context.Context, p *peer.Peer, pmes *Messa
198
return rpmes, nil
199
}
200
200
-func (d *Diagnostics) handleDiagnostic(p *peer.Peer, pmes *Message) (*Message, error) {
201
+func (d *Diagnostics) handleDiagnostic(p peer.Peer, pmes *Message) (*Message, error) {
202
log.Debug("HandleDiagnostic from %s for id = %s", p, pmes.GetDiagID())
203
resp := newMessage(pmes.GetDiagID())
204
d.diagLock.Lock()
@@ -220,7 +221,7 @@ func (d *Diagnostics) handleDiagnostic(p *peer.Peer, pmes *Message) (*Message, e
221
for _, p := range d.getPeers() {
222
log.Debug("Sending diagnostic request to peer: %s", p)
223
sendcount++
223
- go func(p *peer.Peer) {
224
+ go func(p peer.Peer) {
225
out, err := d.getDiagnosticFromPeer(ctx, p, pmes)
226
if err != nil {
227
log.Error("getDiagnostic error: %v", err)
@@ -267,7 +268,7 @@ func (d *Diagnostics) HandleMessage(ctx context.Context, mes msg.NetMessage) msg
268
269
// Print out diagnostic
270
log.Info("[peer: %s] Got message from [%s]\n",
270
- d.self.ID.Pretty(), mPeer.ID.Pretty())
271
+ d.self.ID().Pretty(), mPeer.ID().Pretty())
272
273
// dispatch handler.
274
rpmes, err := d.handleDiagnostic(mPeer, pmes)
exchange/bitswap/bitswap.go
+5
-5
@@ -20,7 +20,7 @@ var log = u.Logger("bitswap")
20
21
// NetMessageSession initializes a BitSwap session that communicates over the
22
// provided NetMessage service
23
-func NetMessageSession(parent context.Context, p *peer.Peer,
23
+func NetMessageSession(parent context.Context, p peer.Peer,
24
net inet.Network, srv inet.Service, directory bsnet.Routing,
25
d ds.Datastore, nice bool) exchange.Interface {
26
@@ -83,7 +83,7 @@ func (bs *bitswap) Block(parent context.Context, k u.Key) (*blocks.Block, error)
83
message.AppendWanted(k)
84
for peerToQuery := range peersToQuery {
85
log.Debug("bitswap got peersToQuery: %s", peerToQuery)
86
- go func(p *peer.Peer) {
86
+ go func(p peer.Peer) {
87
88
log.Debug("bitswap dialing peer: %s", p)
89
err := bs.sender.DialPeer(p)
@@ -131,8 +131,8 @@ func (bs *bitswap) HasBlock(ctx context.Context, blk blocks.Block) error {
131
}
132
133
// TODO(brian): handle errors
134
-func (bs *bitswap) ReceiveMessage(ctx context.Context, p *peer.Peer, incoming bsmsg.BitSwapMessage) (
135
- *peer.Peer, bsmsg.BitSwapMessage) {
134
+func (bs *bitswap) ReceiveMessage(ctx context.Context, p peer.Peer, incoming bsmsg.BitSwapMessage) (
135
+ peer.Peer, bsmsg.BitSwapMessage) {
136
log.Debug("ReceiveMessage from %v", p.Key())
137
138
if p == nil {
@@ -181,7 +181,7 @@ func (bs *bitswap) ReceiveError(err error) {
181
182
// send strives to ensure that accounting is always performed when a message is
183
// sent
184
-func (bs *bitswap) send(ctx context.Context, p *peer.Peer, m bsmsg.BitSwapMessage) {
184
+func (bs *bitswap) send(ctx context.Context, p peer.Peer, m bsmsg.BitSwapMessage) {
185
bs.sender.SendMessage(ctx, p, m)
186
go bs.strategy.MessageSent(p, m)
187
}
exchange/bitswap/bitswap_test.go
+3
-3
@@ -44,7 +44,7 @@ func TestProviderForKeyButNetworkCannotFind(t *testing.T) {
44
g := NewSessionGenerator(net, rs)
45
46
block := blocks.NewBlock([]byte("block"))
47
- rs.Announce(&peer.Peer{}, block.Key()) // but not on network
47
+ rs.Announce(peer.WithIDString("testing"), block.Key()) // but not on network
48
49
solo := g.Next()
50
@@ -263,7 +263,7 @@ func (g *SessionGenerator) Instances(n int) []instance {
263
}
264
265
type instance struct {
266
- peer *peer.Peer
266
+ peer peer.Peer
267
exchange exchange.Interface
268
blockstore bstore.Blockstore
269
}
@@ -274,7 +274,7 @@ type instance struct {
274
// sessions. To safeguard, use the SessionGenerator to generate sessions. It's
275
// just a much better idea.
276
func session(net tn.Network, rs mock.RoutingServer, id peer.ID) instance {
277
- p := &peer.Peer{ID: id}
277
+ p := peer.WithID(id)
278
279
adapter := net.Adapter(p)
280
htc := rs.Client(p)
exchange/bitswap/message/message.go
+2
-2
@@ -19,7 +19,7 @@ type BitSwapMessage interface {
19
20
type Exportable interface {
21
ToProto() *PBMessage
22
- ToNet(p *peer.Peer) (nm.NetMessage, error)
22
+ ToNet(p peer.Peer) (nm.NetMessage, error)
23
}
24
25
// message wraps a proto message for convenience
@@ -82,6 +82,6 @@ func (m *message) ToProto() *PBMessage {
82
return pb
83
}
84
85
-func (m *message) ToNet(p *peer.Peer) (nm.NetMessage, error) {
85
+func (m *message) ToNet(p peer.Peer) (nm.NetMessage, error) {
86
return nm.FromObject(p, m.ToProto())
87
}
exchange/bitswap/message/message_test.go
+4
-3
@@ -88,7 +88,7 @@ func TestCopyProtoByValue(t *testing.T) {
88
89
func TestToNetMethodSetsPeer(t *testing.T) {
90
m := New()
91
- p := &peer.Peer{ID: []byte("X")}
91
+ p := peer.WithIDString("X")
92
netmsg, err := m.ToNet(p)
93
if err != nil {
94
t.Fatal(err)
@@ -106,7 +106,8 @@ func TestToNetFromNetPreservesWantList(t *testing.T) {
106
original.AppendWanted(u.Key("T"))
107
original.AppendWanted(u.Key("F"))
108
109
- netmsg, err := original.ToNet(&peer.Peer{ID: []byte("X")})
109
+ p := peer.WithIDString("X")
110
+ netmsg, err := original.ToNet(p)
111
if err != nil {
112
t.Fatal(err)
113
}
@@ -136,7 +137,7 @@ func TestToAndFromNetMessage(t *testing.T) {
137
original.AppendBlock(*blocks.NewBlock([]byte("F")))
138
original.AppendBlock(*blocks.NewBlock([]byte("M")))
139
139
- p := &peer.Peer{ID: []byte("X")}
140
+ p := peer.WithIDString("X")
141
netmsg, err := original.ToNet(p)
142
if err != nil {
143
t.Fatal(err)
exchange/bitswap/network/interface.go
+6
-6
@@ -12,18 +12,18 @@ import (
12
type Adapter interface {
13
14
// DialPeer ensures there is a connection to peer.
15
- DialPeer(*peer.Peer) error
15
+ DialPeer(peer.Peer) error
16
17
// SendMessage sends a BitSwap message to a peer.
18
SendMessage(
19
context.Context,
20
- *peer.Peer,
20
+ peer.Peer,
21
bsmsg.BitSwapMessage) error
22
23
// SendRequest sends a BitSwap message to a peer and waits for a response.
24
SendRequest(
25
context.Context,
26
- *peer.Peer,
26
+ peer.Peer,
27
bsmsg.BitSwapMessage) (incoming bsmsg.BitSwapMessage, err error)
28
29
// SetDelegate registers the Reciver to handle messages received from the
@@ -33,8 +33,8 @@ type Adapter interface {
33
34
type Receiver interface {
35
ReceiveMessage(
36
- ctx context.Context, sender *peer.Peer, incoming bsmsg.BitSwapMessage) (
37
- destination *peer.Peer, outgoing bsmsg.BitSwapMessage)
36
+ ctx context.Context, sender peer.Peer, incoming bsmsg.BitSwapMessage) (
37
+ destination peer.Peer, outgoing bsmsg.BitSwapMessage)
38
39
ReceiveError(error)
40
}
@@ -42,7 +42,7 @@ type Receiver interface {
42
// TODO rename -> Router?
43
type Routing interface {
44
// FindProvidersAsync returns a channel of providers for the given key
45
- FindProvidersAsync(context.Context, u.Key, int) <-chan *peer.Peer
45
+ FindProvidersAsync(context.Context, u.Key, int) <-chan peer.Peer
46
47
// Provide provides the key to the network
48
Provide(context.Context, u.Key) error
exchange/bitswap/network/net_message_adapter.go
+3
-3
@@ -60,13 +60,13 @@ func (adapter *impl) HandleMessage(
60
return outgoing
61
}
62
63
-func (adapter *impl) DialPeer(p *peer.Peer) error {
63
+func (adapter *impl) DialPeer(p peer.Peer) error {
64
return adapter.net.DialPeer(p)
65
}
66
67
func (adapter *impl) SendMessage(
68
ctx context.Context,
69
- p *peer.Peer,
69
+ p peer.Peer,
70
outgoing bsmsg.BitSwapMessage) error {
71
72
nmsg, err := outgoing.ToNet(p)
@@ -78,7 +78,7 @@ func (adapter *impl) SendMessage(
78
79
func (adapter *impl) SendRequest(
80
ctx context.Context,
81
- p *peer.Peer,
81
+ p peer.Peer,
82
outgoing bsmsg.BitSwapMessage) (bsmsg.BitSwapMessage, error) {
83
84
outgoingMsg, err := outgoing.ToNet(p)
exchange/bitswap/strategy/interface.go
+7
-7
@@ -8,25 +8,25 @@ import (
8
9
type Strategy interface {
10
// Returns a slice of Peers with whom the local node has active sessions
11
- Peers() []*peer.Peer
11
+ Peers() []peer.Peer
12
13
// BlockIsWantedByPeer returns true if peer wants the block given by this
14
// key
15
- BlockIsWantedByPeer(u.Key, *peer.Peer) bool
15
+ BlockIsWantedByPeer(u.Key, peer.Peer) bool
16
17
// ShouldSendTo(Peer) decides whether to send data to this Peer
18
- ShouldSendBlockToPeer(u.Key, *peer.Peer) bool
18
+ ShouldSendBlockToPeer(u.Key, peer.Peer) bool
19
20
// Seed initializes the decider to a deterministic state
21
Seed(int64)
22
23
// MessageReceived records receipt of message for accounting purposes
24
- MessageReceived(*peer.Peer, bsmsg.BitSwapMessage) error
24
+ MessageReceived(peer.Peer, bsmsg.BitSwapMessage) error
25
26
// MessageSent records sending of message for accounting purposes
27
- MessageSent(*peer.Peer, bsmsg.BitSwapMessage) error
27
+ MessageSent(peer.Peer, bsmsg.BitSwapMessage) error
28
29
- NumBytesSentTo(*peer.Peer) uint64
29
+ NumBytesSentTo(peer.Peer) uint64
30
31
- NumBytesReceivedFrom(*peer.Peer) uint64
31
+ NumBytesReceivedFrom(peer.Peer) uint64
32
}
exchange/bitswap/strategy/ledger.go
+2
-2
@@ -12,7 +12,7 @@ import (
12
// access/lookups.
13
type keySet map[u.Key]struct{}
14
15
-func newLedger(p *peer.Peer, strategy strategyFunc) *ledger {
15
+func newLedger(p peer.Peer, strategy strategyFunc) *ledger {
16
return &ledger{
17
wantList: keySet{},
18
Strategy: strategy,
@@ -25,7 +25,7 @@ type ledger struct {
25
lock sync.RWMutex
26
27
// Partner is the remote Peer.
28
- Partner *peer.Peer
28
+ Partner peer.Peer
29
30
// Accounting tracks bytes sent and recieved.
31
Accounting debtRatio
exchange/bitswap/strategy/strategy.go
+9
-9
@@ -37,20 +37,20 @@ type ledgerMap map[peerKey]*ledger
37
type peerKey u.Key
38
39
// Peers returns a list of peers
40
-func (s *strategist) Peers() []*peer.Peer {
41
- response := make([]*peer.Peer, 0)
40
+func (s *strategist) Peers() []peer.Peer {
41
+ response := make([]peer.Peer, 0)
42
for _, ledger := range s.ledgerMap {
43
response = append(response, ledger.Partner)
44
}
45
return response
46
}
47
48
-func (s *strategist) BlockIsWantedByPeer(k u.Key, p *peer.Peer) bool {
48
+func (s *strategist) BlockIsWantedByPeer(k u.Key, p peer.Peer) bool {
49
ledger := s.ledger(p)
50
return ledger.WantListContains(k)
51
}
52
53
-func (s *strategist) ShouldSendBlockToPeer(k u.Key, p *peer.Peer) bool {
53
+func (s *strategist) ShouldSendBlockToPeer(k u.Key, p peer.Peer) bool {
54
ledger := s.ledger(p)
55
return ledger.ShouldSend()
56
}
@@ -59,7 +59,7 @@ func (s *strategist) Seed(int64) {
59
// TODO
60
}
61
62
-func (s *strategist) MessageReceived(p *peer.Peer, m bsmsg.BitSwapMessage) error {
62
+func (s *strategist) MessageReceived(p peer.Peer, m bsmsg.BitSwapMessage) error {
63
// TODO find a more elegant way to handle this check
64
if p == nil {
65
return errors.New("Strategy received nil peer")
@@ -84,7 +84,7 @@ func (s *strategist) MessageReceived(p *peer.Peer, m bsmsg.BitSwapMessage) error
84
// inconsistent. Would need to ensure that Sends and acknowledgement of the
85
// send happen atomically
86
87
-func (s *strategist) MessageSent(p *peer.Peer, m bsmsg.BitSwapMessage) error {
87
+func (s *strategist) MessageSent(p peer.Peer, m bsmsg.BitSwapMessage) error {
88
l := s.ledger(p)
89
for _, block := range m.Blocks() {
90
l.SentBytes(len(block.Data))
@@ -95,16 +95,16 @@ func (s *strategist) MessageSent(p *peer.Peer, m bsmsg.BitSwapMessage) error {
95
return nil
96
}
97
98
-func (s *strategist) NumBytesSentTo(p *peer.Peer) uint64 {
98
+func (s *strategist) NumBytesSentTo(p peer.Peer) uint64 {
99
return s.ledger(p).Accounting.BytesSent
100
}
101
102
-func (s *strategist) NumBytesReceivedFrom(p *peer.Peer) uint64 {
102
+func (s *strategist) NumBytesReceivedFrom(p peer.Peer) uint64 {
103
return s.ledger(p).Accounting.BytesRecv
104
}
105
106
// ledger lazily instantiates a ledger
107
-func (s *strategist) ledger(p *peer.Peer) *ledger {
107
+func (s *strategist) ledger(p peer.Peer) *ledger {
108
l, ok := s.ledgerMap[peerKey(p.Key())]
109
if !ok {
110
l = newLedger(p, s.strategyFunc)
exchange/bitswap/strategy/strategy_test.go
+3
-3
@@ -10,13 +10,13 @@ import (
10
)
11
12
type peerAndStrategist struct {
13
- *peer.Peer
13
+ peer.Peer
14
Strategy
15
}
16
17
func newPeerAndStrategist(idStr string) peerAndStrategist {
18
return peerAndStrategist{
19
- Peer: &peer.Peer{ID: peer.ID(idStr)},
19
+ Peer: peer.WithIDString(idStr),
20
Strategy: New(true),
21
}
22
}
@@ -93,7 +93,7 @@ func TestPeerIsAddedToPeersWhenMessageReceivedOrSent(t *testing.T) {
93
}
94
}
95
96
-func peerIsPartner(p *peer.Peer, s Strategy) bool {
96
+func peerIsPartner(p peer.Peer, s Strategy) bool {
97
for _, partner := range s.Peers() {
98
if partner.Key() == p.Key() {
99
return true
exchange/bitswap/testnet/network.go
+18
-18
@@ -13,20 +13,20 @@ import (
13
)
14
15
type Network interface {
16
- Adapter(*peer.Peer) bsnet.Adapter
16
+ Adapter(peer.Peer) bsnet.Adapter
17
18
- HasPeer(*peer.Peer) bool
18
+ HasPeer(peer.Peer) bool
19
20
SendMessage(
21
ctx context.Context,
22
- from *peer.Peer,
23
- to *peer.Peer,
22
+ from peer.Peer,
23
+ to peer.Peer,
24
message bsmsg.BitSwapMessage) error
25
26
SendRequest(
27
ctx context.Context,
28
- from *peer.Peer,
29
- to *peer.Peer,
28
+ from peer.Peer,
29
+ to peer.Peer,
30
message bsmsg.BitSwapMessage) (
31
incoming bsmsg.BitSwapMessage, err error)
32
}
@@ -43,7 +43,7 @@ type network struct {
43
clients map[util.Key]bsnet.Receiver
44
}
45
46
-func (n *network) Adapter(p *peer.Peer) bsnet.Adapter {
46
+func (n *network) Adapter(p peer.Peer) bsnet.Adapter {
47
client := &networkClient{
48
local: p,
49
network: n,
@@ -52,7 +52,7 @@ func (n *network) Adapter(p *peer.Peer) bsnet.Adapter {
52
return client
53
}
54
55
-func (n *network) HasPeer(p *peer.Peer) bool {
55
+func (n *network) HasPeer(p peer.Peer) bool {
56
_, found := n.clients[p.Key()]
57
return found
58
}
@@ -61,8 +61,8 @@ func (n *network) HasPeer(p *peer.Peer) bool {
61
// TODO what does the network layer do with errors received from services?
62
func (n *network) SendMessage(
63
ctx context.Context,
64
- from *peer.Peer,
65
- to *peer.Peer,
64
+ from peer.Peer,
65
+ to peer.Peer,
66
message bsmsg.BitSwapMessage) error {
67
68
receiver, ok := n.clients[to.Key()]
@@ -79,7 +79,7 @@ func (n *network) SendMessage(
79
}
80
81
func (n *network) deliver(
82
- r bsnet.Receiver, from *peer.Peer, message bsmsg.BitSwapMessage) error {
82
+ r bsnet.Receiver, from peer.Peer, message bsmsg.BitSwapMessage) error {
83
if message == nil || from == nil {
84
return errors.New("Invalid input")
85
}
@@ -107,8 +107,8 @@ var NoResponse = errors.New("No response received from the receiver")
107
// TODO
108
func (n *network) SendRequest(
109
ctx context.Context,
110
- from *peer.Peer,
111
- to *peer.Peer,
110
+ from peer.Peer,
111
+ to peer.Peer,
112
message bsmsg.BitSwapMessage) (
113
incoming bsmsg.BitSwapMessage, err error) {
114
@@ -130,7 +130,7 @@ func (n *network) SendRequest(
130
}
131
132
// TODO test when receiver doesn't immediately respond to the initiator of the request
133
- if !bytes.Equal(nextPeer.ID, from.ID) {
133
+ if !bytes.Equal(nextPeer.ID(), from.ID()) {
134
go func() {
135
nextReceiver, ok := n.clients[nextPeer.Key()]
136
if !ok {
@@ -144,26 +144,26 @@ func (n *network) SendRequest(
144
}
145
146
type networkClient struct {
147
- local *peer.Peer
147
+ local peer.Peer
148
bsnet.Receiver
149
network Network
150
}
151
152
func (nc *networkClient) SendMessage(
153
ctx context.Context,
154
- to *peer.Peer,
154
+ to peer.Peer,
155
message bsmsg.BitSwapMessage) error {
156
return nc.network.SendMessage(ctx, nc.local, to, message)
157
}
158
159
func (nc *networkClient) SendRequest(
160
ctx context.Context,
161
- to *peer.Peer,
161
+ to peer.Peer,
162
message bsmsg.BitSwapMessage) (incoming bsmsg.BitSwapMessage, err error) {
163
return nc.network.SendRequest(ctx, nc.local, to, message)
164
}
165
166
-func (nc *networkClient) DialPeer(p *peer.Peer) error {
166
+func (nc *networkClient) DialPeer(p peer.Peer) error {
167
// no need to do anything because dialing isn't a thing in this test net.
168
if !nc.network.HasPeer(p) {
169
return fmt.Errorf("Peer not in network: %s", p)
exchange/bitswap/testnet/network_test.go
+18
-18
@@ -18,15 +18,15 @@ func TestSendRequestToCooperativePeer(t *testing.T) {
18
19
t.Log("Get two network adapters")
20
21
- initiator := net.Adapter(&peer.Peer{ID: []byte("initiator")})
22
- recipient := net.Adapter(&peer.Peer{ID: idOfRecipient})
21
+ initiator := net.Adapter(peer.WithIDString("initiator"))
22
+ recipient := net.Adapter(peer.WithID(idOfRecipient))
23
24
expectedStr := "response from recipient"
25
recipient.SetDelegate(lambda(func(
26
ctx context.Context,
27
- from *peer.Peer,
27
+ from peer.Peer,
28
incoming bsmsg.BitSwapMessage) (
29
- *peer.Peer, bsmsg.BitSwapMessage) {
29
+ peer.Peer, bsmsg.BitSwapMessage) {
30
31
t.Log("Recipient received a message from the network")
32
@@ -43,7 +43,7 @@ func TestSendRequestToCooperativePeer(t *testing.T) {
43
message := bsmsg.New()
44
message.AppendBlock(*blocks.NewBlock([]byte("data")))
45
response, err := initiator.SendRequest(
46
- context.Background(), &peer.Peer{ID: idOfRecipient}, message)
46
+ context.Background(), peer.WithID(idOfRecipient), message)
47
if err != nil {
48
t.Fatal(err)
49
}
@@ -61,8 +61,8 @@ func TestSendRequestToCooperativePeer(t *testing.T) {
61
func TestSendMessageAsyncButWaitForResponse(t *testing.T) {
62
net := VirtualNetwork()
63
idOfResponder := []byte("responder")
64
- waiter := net.Adapter(&peer.Peer{ID: []byte("waiter")})
65
- responder := net.Adapter(&peer.Peer{ID: idOfResponder})
64
+ waiter := net.Adapter(peer.WithIDString("waiter"))
65
+ responder := net.Adapter(peer.WithID(idOfResponder))
66
67
var wg sync.WaitGroup
68
@@ -72,9 +72,9 @@ func TestSendMessageAsyncButWaitForResponse(t *testing.T) {
72
73
responder.SetDelegate(lambda(func(
74
ctx context.Context,
75
- fromWaiter *peer.Peer,
75
+ fromWaiter peer.Peer,
76
msgFromWaiter bsmsg.BitSwapMessage) (
77
- *peer.Peer, bsmsg.BitSwapMessage) {
77
+ peer.Peer, bsmsg.BitSwapMessage) {
78
79
msgToWaiter := bsmsg.New()
80
msgToWaiter.AppendBlock(*blocks.NewBlock([]byte(expectedStr)))
@@ -84,9 +84,9 @@ func TestSendMessageAsyncButWaitForResponse(t *testing.T) {
84
85
waiter.SetDelegate(lambda(func(
86
ctx context.Context,
87
- fromResponder *peer.Peer,
87
+ fromResponder peer.Peer,
88
msgFromResponder bsmsg.BitSwapMessage) (
89
- *peer.Peer, bsmsg.BitSwapMessage) {
89
+ peer.Peer, bsmsg.BitSwapMessage) {
90
91
// TODO assert that this came from the correct peer and that the message contents are as expected
92
ok := false
@@ -107,7 +107,7 @@ func TestSendMessageAsyncButWaitForResponse(t *testing.T) {
107
messageSentAsync := bsmsg.New()
108
messageSentAsync.AppendBlock(*blocks.NewBlock([]byte("data")))
109
errSending := waiter.SendMessage(
110
- context.Background(), &peer.Peer{ID: idOfResponder}, messageSentAsync)
110
+ context.Background(), peer.WithID(idOfResponder), messageSentAsync)
111
if errSending != nil {
112
t.Fatal(errSending)
113
}
@@ -115,8 +115,8 @@ func TestSendMessageAsyncButWaitForResponse(t *testing.T) {
115
wg.Wait() // until waiter delegate function is executed
116
}
117
118
-type receiverFunc func(ctx context.Context, p *peer.Peer,
119
- incoming bsmsg.BitSwapMessage) (*peer.Peer, bsmsg.BitSwapMessage)
118
+type receiverFunc func(ctx context.Context, p peer.Peer,
119
+ incoming bsmsg.BitSwapMessage) (peer.Peer, bsmsg.BitSwapMessage)
120
121
// lambda returns a Receiver instance given a receiver function
122
func lambda(f receiverFunc) bsnet.Receiver {
@@ -126,13 +126,13 @@ func lambda(f receiverFunc) bsnet.Receiver {
126
}
127
128
type lambdaImpl struct {
129
- f func(ctx context.Context, p *peer.Peer, incoming bsmsg.BitSwapMessage) (
130
- *peer.Peer, bsmsg.BitSwapMessage)
129
+ f func(ctx context.Context, p peer.Peer, incoming bsmsg.BitSwapMessage) (
130
+ peer.Peer, bsmsg.BitSwapMessage)
131
}
132
133
func (lam *lambdaImpl) ReceiveMessage(ctx context.Context,
134
- p *peer.Peer, incoming bsmsg.BitSwapMessage) (
135
- *peer.Peer, bsmsg.BitSwapMessage) {
134
+ p peer.Peer, incoming bsmsg.BitSwapMessage) (
135
+ peer.Peer, bsmsg.BitSwapMessage) {
136
return lam.f(ctx, p, incoming)
137
}
138
fuse/ipns/ipns_test.go
+1
-1
@@ -209,7 +209,7 @@ func TestFastRepublish(t *testing.T) {
209
210
node, mnt := setupIpnsTest(t, nil)
211
212
- h, err := node.Identity.PrivKey.GetPublic().Hash()
212
+ h, err := node.Identity.PrivKey().GetPublic().Hash()
213
if err != nil {
214
t.Fatal(err)
215
}
fuse/ipns/ipns_unix.go
+1
-1
@@ -35,7 +35,7 @@ type FileSystem struct {
35
36
// NewFileSystem constructs new fs using given core.IpfsNode instance.
37
func NewIpns(ipfs *core.IpfsNode, ipfspath string) (*FileSystem, error) {
38
- root, err := CreateRoot(ipfs, []ci.PrivKey{ipfs.Identity.PrivKey}, ipfspath)
38
+ root, err := CreateRoot(ipfs, []ci.PrivKey{ipfs.Identity.PrivKey()}, ipfspath)
39
if err != nil {
40
return nil, err
41
}
namesys/resolve_test.go
+1
-3
@@ -11,9 +11,7 @@ import (
11
)
12
13
func TestRoutingResolve(t *testing.T) {
14
- local := &peer.Peer{
15
- ID: []byte("testID"),
16
- }
14
+ local := peer.WithIDString("testID")
15
lds := ds.NewMapDatastore()
16
d := mock.NewMockRouter(local, lds)
17
net/conn/conn.go
+5
-5
@@ -42,8 +42,8 @@ func newMsgioPipe(size int) *msgioPipe {
42
43
// singleConn represents a single connection to another Peer (IPFS Node).
44
type singleConn struct {
45
- local *peer.Peer
46
- remote *peer.Peer
45
+ local peer.Peer
46
+ remote peer.Peer
47
maconn manet.Conn
48
msgio *msgioPipe
49
@@ -51,7 +51,7 @@ type singleConn struct {
51
}
52
53
// newConn constructs a new connection
54
-func newSingleConn(ctx context.Context, local, remote *peer.Peer,
54
+func newSingleConn(ctx context.Context, local, remote peer.Peer,
55
maconn manet.Conn) (Conn, error) {
56
57
conn := &singleConn{
@@ -117,12 +117,12 @@ func (c *singleConn) RemoteMultiaddr() ma.Multiaddr {
117
}
118
119
// LocalPeer is the Peer on this side
120
-func (c *singleConn) LocalPeer() *peer.Peer {
120
+func (c *singleConn) LocalPeer() peer.Peer {
121
return c.local
122
}
123
124
// RemotePeer is the Peer on the remote side
125
-func (c *singleConn) RemotePeer() *peer.Peer {
125
+func (c *singleConn) RemotePeer() peer.Peer {
126
return c.remote
127
}
128
net/conn/dial.go
+1
-1
@@ -12,7 +12,7 @@ import (
12
13
// Dial connects to a particular peer, over a given network
14
// Example: d.Dial(ctx, "udp", peer)
15
-func (d *Dialer) Dial(ctx context.Context, network string, remote *peer.Peer) (Conn, error) {
15
+func (d *Dialer) Dial(ctx context.Context, network string, remote peer.Peer) (Conn, error) {
16
laddr := d.LocalPeer.NetAddress(network)
17
if laddr == nil {
18
return nil, fmt.Errorf("No local address for network %s", network)
net/conn/dial_test.go
+2
-6
@@ -10,7 +10,7 @@ import (
10
ma "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-multiaddr"
11
)
12
13
-func setupPeer(addr string) (*peer.Peer, error) {
13
+func setupPeer(addr string) (peer.Peer, error) {
14
tcp, err := ma.NewMultiaddr(addr)
15
if err != nil {
16
return nil, err
@@ -21,14 +21,10 @@ func setupPeer(addr string) (*peer.Peer, error) {
21
return nil, err
22
}
23
24
- id, err := peer.IDFromPubKey(pk)
24
+ p, err := peer.WithKeyPair(sk, pk)
25
if err != nil {
26
return nil, err
27
}
28
-
29
- p := &peer.Peer{ID: id}
30
- p.PrivKey = sk
31
- p.PubKey = pk
28
p.AddAddress(tcp)
29
return p, nil
30
}
net/conn/interface.go
+4
-4
@@ -23,13 +23,13 @@ type Conn interface {
23
LocalMultiaddr() ma.Multiaddr
24
25
// LocalPeer is the Peer on this side
26
- LocalPeer() *peer.Peer
26
+ LocalPeer() peer.Peer
27
28
// RemoteMultiaddr is the Multiaddr on the remote side
29
RemoteMultiaddr() ma.Multiaddr
30
31
// RemotePeer is the Peer on the remote side
32
- RemotePeer() *peer.Peer
32
+ RemotePeer() peer.Peer
33
34
// In returns a readable message channel
35
In() <-chan []byte
@@ -47,7 +47,7 @@ type Conn interface {
47
type Dialer struct {
48
49
// LocalPeer is the identity of the local Peer.
50
- LocalPeer *peer.Peer
50
+ LocalPeer peer.Peer
51
52
// Peerstore is the set of peers we know about locally. The Dialer needs it
53
// because when an incoming connection is identified, we should reuse the
@@ -65,7 +65,7 @@ type Listener interface {
65
Multiaddr() ma.Multiaddr
66
67
// LocalPeer is the identity of the local Peer.
68
- LocalPeer() *peer.Peer
68
+ LocalPeer() peer.Peer
69
70
// Peerstore is the set of peers we know about locally. The Listener needs it
71
// because when an incoming connection is identified, we should reuse the
net/conn/listen.go
+3
-3
@@ -23,7 +23,7 @@ type listener struct {
23
maddr ma.Multiaddr
24
25
// LocalPeer is the identity of the local Peer.
26
- local *peer.Peer
26
+ local peer.Peer
27
28
// Peerstore is the set of peers we know about locally
29
peers peer.Peerstore
@@ -105,7 +105,7 @@ func (l *listener) Multiaddr() ma.Multiaddr {
105
}
106
107
// LocalPeer is the identity of the local Peer.
108
-func (l *listener) LocalPeer() *peer.Peer {
108
+func (l *listener) LocalPeer() peer.Peer {
109
return l.local
110
}
111
@@ -117,7 +117,7 @@ func (l *listener) Peerstore() peer.Peerstore {
117
}
118
119
// Listen listens on the particular multiaddr, with given peer and peerstore.
120
-func Listen(ctx context.Context, addr ma.Multiaddr, local *peer.Peer, peers peer.Peerstore) (Listener, error) {
120
+func Listen(ctx context.Context, addr ma.Multiaddr, local peer.Peer, peers peer.Peerstore) (Listener, error) {
121
122
ml, err := manet.Listen(addr)
123
if err != nil {
net/conn/multiconn.go
+9
-5
@@ -27,8 +27,8 @@ type MultiConn struct {
27
// this string is: /addr1/peer1/addr2/peer2 (peers ordered lexicographically)
28
conns map[string]Conn
29
30
- local *peer.Peer
31
- remote *peer.Peer
30
+ local peer.Peer
31
+ remote peer.Peer
32
33
// fan-in/fan-out
34
duplex Duplex
@@ -39,7 +39,7 @@ type MultiConn struct {
39
}
40
41
// NewMultiConn constructs a new connection
42
-func NewMultiConn(ctx context.Context, local, remote *peer.Peer, conns []Conn) (*MultiConn, error) {
42
+func NewMultiConn(ctx context.Context, local, remote peer.Peer, conns []Conn) (*MultiConn, error) {
43
44
c := &MultiConn{
45
local: local,
@@ -72,6 +72,10 @@ func (c *MultiConn) Add(conns ...Conn) {
72
log.Error("%s", c2)
73
c.Unlock() // ok to unlock (to log). panicing.
74
log.Error("%s", c)
75
+ log.Error("c.LocalPeer: %s %#v", c.LocalPeer(), c.LocalPeer())
76
+ log.Error("c2.LocalPeer: %s %#v", c2.LocalPeer(), c2.LocalPeer())
77
+ log.Error("c.RemotePeer: %s %#v", c.RemotePeer(), c.RemotePeer())
78
+ log.Error("c2.RemotePeer: %s %#v", c2.RemotePeer(), c2.RemotePeer())
79
c.Lock() // gotta relock to avoid lock panic from deferring.
80
panic("connection addresses mismatch")
81
}
@@ -269,12 +273,12 @@ func (c *MultiConn) RemoteMultiaddr() ma.Multiaddr {
273
}
274
275
// LocalPeer is the Peer on this side
272
-func (c *MultiConn) LocalPeer() *peer.Peer {
276
+func (c *MultiConn) LocalPeer() peer.Peer {
277
return c.local
278
}
279
280
// RemotePeer is the Peer on the remote side
277
-func (c *MultiConn) RemotePeer() *peer.Peer {
281
+func (c *MultiConn) RemotePeer() peer.Peer {
282
return c.remote
283
}
284
net/conn/multiconn_test.go
+4
-4
@@ -97,7 +97,7 @@ func setupMultiConns(t *testing.T, ctx context.Context) (a, b *MultiConn) {
97
p2ps := peer.NewPeerstore()
98
99
// listeners
100
- listen := func(addr ma.Multiaddr, p *peer.Peer, ps peer.Peerstore) Listener {
100
+ listen := func(addr ma.Multiaddr, p peer.Peer, ps peer.Peerstore) Listener {
101
l, err := Listen(ctx, addr, p, ps)
102
if err != nil {
103
t.Fatal(err)
@@ -106,14 +106,14 @@ func setupMultiConns(t *testing.T, ctx context.Context) (a, b *MultiConn) {
106
}
107
108
log.Info("Setting up listeners")
109
- p1l := listen(p1.Addresses[0], p1, p1ps)
110
- p2l := listen(p2.Addresses[0], p2, p2ps)
109
+ p1l := listen(p1.Addresses()[0], p1, p1ps)
110
+ p2l := listen(p2.Addresses()[0], p2, p2ps)
111
112
// dialers
113
p1d := &Dialer{Peerstore: p1ps, LocalPeer: p1}
114
p2d := &Dialer{Peerstore: p2ps, LocalPeer: p2}
115
116
- dial := func(d *Dialer, dst *peer.Peer) <-chan Conn {
116
+ dial := func(d *Dialer, dst peer.Peer) <-chan Conn {
117
cc := make(chan Conn)
118
go func() {
119
c, err := d.Dial(ctx, "tcp", dst)
net/conn/secure_conn.go
+4
-2
@@ -79,6 +79,8 @@ func (c *secureConn) secureHandshake(peers peer.Peerstore) error {
79
// perhaps return an error. TBD.
80
81
log.Error("secureConn peer mismatch. %v != %v", insecureSC.remote, c.secure.RemotePeer())
82
+ log.Error("insecureSC.remote: %s %#v", insecureSC.remote, insecureSC.remote)
83
+ log.Error("c.secure.LocalPeer: %s %#v", c.secure.RemotePeer(), c.secure.RemotePeer())
84
panic("secureConn peer mismatch. consructed incorrectly?")
85
}
86
@@ -114,12 +116,12 @@ func (c *secureConn) RemoteMultiaddr() ma.Multiaddr {
116
}
117
118
// LocalPeer is the Peer on this side
117
-func (c *secureConn) LocalPeer() *peer.Peer {
119
+func (c *secureConn) LocalPeer() peer.Peer {
120
return c.insecure.LocalPeer()
121
}
122
123
// RemotePeer is the Peer on the remote side
122
-func (c *secureConn) RemotePeer() *peer.Peer {
124
+func (c *secureConn) RemotePeer() peer.Peer {
125
return c.insecure.RemotePeer()
126
}
127
net/interface.go
+4
-4
@@ -15,19 +15,19 @@ type Network interface {
15
// TODO: for now, only listen on addrs in local peer when initializing.
16
17
// DialPeer attempts to establish a connection to a given peer
18
- DialPeer(*peer.Peer) error
18
+ DialPeer(peer.Peer) error
19
20
// ClosePeer connection to peer
21
- ClosePeer(*peer.Peer) error
21
+ ClosePeer(peer.Peer) error
22
23
// IsConnected returns whether a connection to given peer exists.
24
- IsConnected(*peer.Peer) (bool, error)
24
+ IsConnected(peer.Peer) (bool, error)
25
26
// GetProtocols returns the protocols registered in the network.
27
GetProtocols() *mux.ProtocolMap
28
29
// GetPeerList returns the list of peers currently connected in this network.
30
- GetPeerList() []*peer.Peer
30
+ GetPeerList() []peer.Peer
31
32
// GetBandwidthTotals returns the total number of bytes passed through
33
// the network since it was instantiated
net/message/message.go
+5
-5
@@ -8,12 +8,12 @@ import (
8
9
// NetMessage is the interface for the message
10
type NetMessage interface {
11
- Peer() *peer.Peer
11
+ Peer() peer.Peer
12
Data() []byte
13
}
14
15
// New is the interface for constructing a new message.
16
-func New(p *peer.Peer, data []byte) NetMessage {
16
+func New(p peer.Peer, data []byte) NetMessage {
17
return &message{peer: p, data: data}
18
}
19
@@ -21,13 +21,13 @@ func New(p *peer.Peer, data []byte) NetMessage {
21
// particular Peer.
22
type message struct {
23
// To or from, depending on direction.
24
- peer *peer.Peer
24
+ peer peer.Peer
25
26
// Opaque data
27
data []byte
28
}
29
30
-func (m *message) Peer() *peer.Peer {
30
+func (m *message) Peer() peer.Peer {
31
return m.peer
32
}
33
@@ -36,7 +36,7 @@ func (m *message) Data() []byte {
36
}
37
38
// FromObject creates a message from a protobuf-marshallable message.
39
-func FromObject(p *peer.Peer, data proto.Message) (NetMessage, error) {
39
+func FromObject(p peer.Peer, data proto.Message) (NetMessage, error) {
40
bytes, err := proto.Marshal(data)
41
if err != nil {
42
return nil, err
net/mux/mux_test.go
+2
-2
@@ -21,14 +21,14 @@ func (t *TestProtocol) GetPipe() *msg.Pipe {
21
return t.Pipe
22
}
23
24
-func newPeer(t *testing.T, id string) *peer.Peer {
24
+func newPeer(t *testing.T, id string) peer.Peer {
25
mh, err := mh.FromHexString(id)
26
if err != nil {
27
t.Error(err)
28
return nil
29
}
30
31
- return &peer.Peer{ID: peer.ID(mh)}
31
+ return peer.WithID(peer.ID(mh))
32
}
33
34
func testMsg(t *testing.T, m msg.NetMessage, data []byte) {
net/net.go
+9
-7
@@ -15,7 +15,7 @@ import (
15
type IpfsNetwork struct {
16
17
// local peer
18
- local *peer.Peer
18
+ local peer.Peer
19
20
// protocol multiplexing
21
muxer *mux.Muxer
@@ -29,7 +29,7 @@ type IpfsNetwork struct {
29
}
30
31
// NewIpfsNetwork is the structure that implements the network interface
32
-func NewIpfsNetwork(ctx context.Context, local *peer.Peer,
32
+func NewIpfsNetwork(ctx context.Context, local peer.Peer,
33
peers peer.Peerstore, pmap *mux.ProtocolMap) (*IpfsNetwork, error) {
34
35
ctx, cancel := context.WithCancel(ctx)
@@ -63,19 +63,19 @@ func NewIpfsNetwork(ctx context.Context, local *peer.Peer,
63
// func (n *IpfsNetwork) Listen(*ma.Muliaddr) error {}
64
65
// DialPeer attempts to establish a connection to a given peer
66
-func (n *IpfsNetwork) DialPeer(p *peer.Peer) error {
66
+func (n *IpfsNetwork) DialPeer(p peer.Peer) error {
67
_, err := n.swarm.Dial(p)
68
return err
69
}
70
71
// ClosePeer connection to peer
72
-func (n *IpfsNetwork) ClosePeer(p *peer.Peer) error {
72
+func (n *IpfsNetwork) ClosePeer(p peer.Peer) error {
73
return n.swarm.CloseConnection(p)
74
}
75
76
// IsConnected returns whether a connection to given peer exists.
77
-func (n *IpfsNetwork) IsConnected(p *peer.Peer) (bool, error) {
78
- return n.swarm.GetConnection(p.ID) != nil, nil
77
+func (n *IpfsNetwork) IsConnected(p peer.Peer) (bool, error) {
78
+ return n.swarm.GetConnection(p.ID()) != nil, nil
79
}
80
81
// GetProtocols returns the protocols registered in the network.
@@ -108,10 +108,12 @@ func (n *IpfsNetwork) Close() error {
108
return nil
109
}
110
111
-func (n *IpfsNetwork) GetPeerList() []*peer.Peer {
111
+// GetPeerList returns the networks list of connected peers
112
+func (n *IpfsNetwork) GetPeerList() []peer.Peer {
113
return n.swarm.GetPeerList()
114
}
115
116
+// GetBandwidthTotals returns the total amount of bandwidth transferred
117
func (n *IpfsNetwork) GetBandwidthTotals() (in uint64, out uint64) {
118
return n.muxer.GetBandwidthTotals()
119
}
net/service/service.go
+2
-2
@@ -133,7 +133,7 @@ func (s *service) SendMessage(ctx context.Context, m msg.NetMessage) error {
133
func (s *service) SendRequest(ctx context.Context, m msg.NetMessage) (msg.NetMessage, error) {
134
135
// create a request
136
- r, err := NewRequest(m.Peer().ID)
136
+ r, err := NewRequest(m.Peer().ID())
137
if err != nil {
138
return nil, err
139
}
@@ -227,7 +227,7 @@ func (s *service) handleIncomingMessage(ctx context.Context, m msg.NetMessage) {
227
log.Error("RequestID should identify a response here.")
228
}
229
230
- key := RequestKey(m.Peer().ID, RequestID(rid))
230
+ key := RequestKey(m.Peer().ID(), RequestID(rid))
231
s.RequestsLock.RLock()
232
r, found := s.Requests[key]
233
s.RequestsLock.RUnlock()
net/service/service_test.go
+2
-2
@@ -25,14 +25,14 @@ func (t *ReverseHandler) HandleMessage(ctx context.Context, m msg.NetMessage) ms
25
return msg.New(m.Peer(), d)
26
}
27
28
-func newPeer(t *testing.T, id string) *peer.Peer {
28
+func newPeer(t *testing.T, id string) peer.Peer {
29
mh, err := mh.FromHexString(id)
30
if err != nil {
31
t.Error(err)
32
return nil
33
}
34
35
- return &peer.Peer{ID: peer.ID(mh)}
35
+ return peer.WithID(peer.ID(mh))
36
}
37
38
func TestServiceHandler(t *testing.T) {
net/swarm/conn.go
+2
-2
@@ -14,11 +14,11 @@ import (
14
func (s *Swarm) listen() error {
15
hasErr := false
16
retErr := &ListenErr{
17
- Errors: make([]error, len(s.local.Addresses)),
17
+ Errors: make([]error, len(s.local.Addresses())),
18
}
19
20
// listen on every address
21
- for i, addr := range s.local.Addresses {
21
+ for i, addr := range s.local.Addresses() {
22
err := s.connListen(addr)
23
if err != nil {
24
hasErr = true
net/swarm/simul_test.go
+3
-3
@@ -24,10 +24,10 @@ func TestSimultOpen(t *testing.T) {
24
// connect everyone
25
{
26
var wg sync.WaitGroup
27
- connect := func(s *Swarm, dst *peer.Peer) {
27
+ connect := func(s *Swarm, dst peer.Peer) {
28
// copy for other peer
29
- cp := &peer.Peer{ID: dst.ID}
30
- cp.AddAddress(dst.Addresses[0])
29
+ cp := peer.WithID(dst.ID())
30
+ cp.AddAddress(dst.Addresses()[0])
31
32
if _, err := s.Dial(cp); err != nil {
33
t.Fatal("error swarm dialing to peer", err)
net/swarm/swarm.go
+10
-10
@@ -45,7 +45,7 @@ func (e *ListenErr) Error() string {
45
type Swarm struct {
46
47
// local is the peer this swarm represents
48
- local *peer.Peer
48
+ local peer.Peer
49
50
// peers is a collection of peers for swarm to use
51
peers peer.Peerstore
@@ -69,7 +69,7 @@ type Swarm struct {
69
}
70
71
// NewSwarm constructs a Swarm, with a Chan.
72
-func NewSwarm(ctx context.Context, local *peer.Peer, ps peer.Peerstore) (*Swarm, error) {
72
+func NewSwarm(ctx context.Context, local peer.Peer, ps peer.Peerstore) (*Swarm, error) {
73
s := &Swarm{
74
Pipe: msg.NewPipe(10),
75
conns: conn.MultiConnMap{},
@@ -104,13 +104,13 @@ func (s *Swarm) close() error {
104
// etc. to achive connection.
105
//
106
// For now, Dial uses only TCP. This will be extended.
107
-func (s *Swarm) Dial(peer *peer.Peer) (conn.Conn, error) {
108
- if peer.ID.Equal(s.local.ID) {
107
+func (s *Swarm) Dial(peer peer.Peer) (conn.Conn, error) {
108
+ if peer.ID().Equal(s.local.ID()) {
109
return nil, errors.New("Attempted connection to self!")
110
}
111
112
// check if we already have an open connection first
113
- c := s.GetConnection(peer.ID)
113
+ c := s.GetConnection(peer.ID())
114
if c != nil {
115
return c, nil
116
}
@@ -167,14 +167,14 @@ func (s *Swarm) Connections() []conn.Conn {
167
}
168
169
// CloseConnection removes a given peer from swarm + closes the connection
170
-func (s *Swarm) CloseConnection(p *peer.Peer) error {
171
- c := s.GetConnection(p.ID)
170
+func (s *Swarm) CloseConnection(p peer.Peer) error {
171
+ c := s.GetConnection(p.ID())
172
if c == nil {
173
return u.ErrNotFound
174
}
175
176
s.connsLock.Lock()
177
- delete(s.conns, u.Key(p.ID))
177
+ delete(s.conns, u.Key(p.ID()))
178
s.connsLock.Unlock()
179
180
return c.Close()
@@ -190,8 +190,8 @@ func (s *Swarm) GetErrChan() chan error {
190
}
191
192
// GetPeerList returns a copy of the set of peers swarm is connected to.
193
-func (s *Swarm) GetPeerList() []*peer.Peer {
194
- var out []*peer.Peer
193
+func (s *Swarm) GetPeerList() []peer.Peer {
194
+ var out []peer.Peer
195
s.connsLock.RLock()
196
for _, p := range s.conns {
197
out = append(out, p.RemotePeer())
net/swarm/swarm_test.go
+8
-12
@@ -32,7 +32,7 @@ func pong(ctx context.Context, swarm *Swarm) {
32
}
33
}
34
35
-func setupPeer(t *testing.T, addr string) *peer.Peer {
35
+func setupPeer(t *testing.T, addr string) peer.Peer {
36
tcp, err := ma.NewMultiaddr(addr)
37
if err != nil {
38
t.Fatal(err)
@@ -43,19 +43,15 @@ func setupPeer(t *testing.T, addr string) *peer.Peer {
43
t.Fatal(err)
44
}
45
46
- id, err := peer.IDFromPubKey(pk)
46
+ p, err := peer.WithKeyPair(sk, pk)
47
if err != nil {
48
t.Fatal(err)
49
}
50
-
51
- p := &peer.Peer{ID: id}
52
- p.PrivKey = sk
53
- p.PubKey = pk
50
p.AddAddress(tcp)
51
return p
52
}
53
58
-func makeSwarms(ctx context.Context, t *testing.T, addrs []string) ([]*Swarm, []*peer.Peer) {
54
+func makeSwarms(ctx context.Context, t *testing.T, addrs []string) ([]*Swarm, []peer.Peer) {
55
swarms := []*Swarm{}
56
57
for _, addr := range addrs {
@@ -68,7 +64,7 @@ func makeSwarms(ctx context.Context, t *testing.T, addrs []string) ([]*Swarm, []
64
swarms = append(swarms, swarm)
65
}
66
71
- peers := make([]*peer.Peer, len(swarms))
67
+ peers := make([]peer.Peer, len(swarms))
68
for i, s := range swarms {
69
peers[i] = s.local
70
}
@@ -85,14 +81,14 @@ func SubtestSwarm(t *testing.T, addrs []string, MsgNum int) {
81
// connect everyone
82
{
83
var wg sync.WaitGroup
88
- connect := func(s *Swarm, dst *peer.Peer) {
84
+ connect := func(s *Swarm, dst peer.Peer) {
85
// copy for other peer
86
91
- cp, err := s.peers.Get(dst.ID)
87
+ cp, err := s.peers.Get(dst.ID())
88
if err != nil {
93
- cp = &peer.Peer{ID: dst.ID}
89
+ t.Fatal(err)
90
}
95
- cp.AddAddress(dst.Addresses[0])
91
+ cp.AddAddress(dst.Addresses()[0])
92
93
log.Info("SWARM TEST: %s dialing %s", s.local, dst)
94
if _, err := s.Dial(cp); err != nil {
peer/peer.go
+119
-29
@@ -49,17 +49,48 @@ func IDFromPubKey(pk ic.PubKey) (ID, error) {
49
return ID(hash), nil
50
}
51
52
-// Map maps Key (string) : *Peer (slices are not comparable).
53
-type Map map[u.Key]*Peer
52
+// Map maps Key (string) : *peer (slices are not comparable).
53
+type Map map[u.Key]Peer
54
55
// Peer represents the identity information of an IPFS Node, including
56
// ID, and relevant Addresses.
57
-type Peer struct {
58
- ID ID
59
- Addresses []ma.Multiaddr
57
+type Peer interface {
58
+ // ID returns the peer's ID
59
+ ID() ID
60
61
- PrivKey ic.PrivKey
62
- PubKey ic.PubKey
61
+ // Key returns the ID as a Key (string) for maps.
62
+ Key() u.Key
63
+
64
+ // Addresses returns the peer's multiaddrs
65
+ Addresses() []ma.Multiaddr
66
+
67
+ // AddAddress adds the given Multiaddr address to Peer's addresses.
68
+ AddAddress(a ma.Multiaddr)
69
+
70
+ // NetAddress returns the first Multiaddr found for a given network.
71
+ NetAddress(n string) ma.Multiaddr
72
+
73
+ // Priv/PubKey returns the peer's Private Key
74
+ PrivKey() ic.PrivKey
75
+ PubKey() ic.PubKey
76
+
77
+ // LoadAndVerifyKeyPair unmarshalls, loads a private/public key pair.
78
+ // Error if (a) unmarshalling fails, or (b) pubkey does not match id.
79
+ LoadAndVerifyKeyPair(marshalled []byte) error
80
+ VerifyAndSetPrivKey(sk ic.PrivKey) error
81
+ VerifyAndSetPubKey(pk ic.PubKey) error
82
+
83
+ // Get/SetLatency manipulate the current latency measurement.
84
+ GetLatency() (out time.Duration)
85
+ SetLatency(laten time.Duration)
86
+}
87
+
88
+type peer struct {
89
+ id ID
90
+ addresses []ma.Multiaddr
91
+
92
+ privKey ic.PrivKey
93
+ pubKey ic.PubKey
94
95
latency time.Duration
96
@@ -67,34 +98,56 @@ type Peer struct {
98
}
99
100
// String prints out the peer.
70
-func (p *Peer) String() string {
71
- return "[Peer " + p.ID.String()[:12] + "]"
101
+func (p *peer) String() string {
102
+ return "[Peer " + p.id.String()[:12] + "]"
103
}
104
105
// Key returns the ID as a Key (string) for maps.
75
-func (p *Peer) Key() u.Key {
76
- return u.Key(p.ID)
106
+func (p *peer) Key() u.Key {
107
+ return u.Key(p.id)
108
+}
109
+
110
+// ID returns the peer's ID
111
+func (p *peer) ID() ID {
112
+ return p.id
113
+}
114
+
115
+// PrivKey returns the peer's Private Key
116
+func (p *peer) PrivKey() ic.PrivKey {
117
+ return p.privKey
118
+}
119
+
120
+// PubKey returns the peer's Private Key
121
+func (p *peer) PubKey() ic.PubKey {
122
+ return p.pubKey
123
+}
124
+
125
+// Addresses returns the peer's multiaddrs
126
+func (p *peer) Addresses() []ma.Multiaddr {
127
+ cp := make([]ma.Multiaddr, len(p.addresses))
128
+ copy(cp, p.addresses)
129
+ return cp
130
}
131
132
// AddAddress adds the given Multiaddr address to Peer's addresses.
80
-func (p *Peer) AddAddress(a ma.Multiaddr) {
133
+func (p *peer) AddAddress(a ma.Multiaddr) {
134
p.Lock()
135
defer p.Unlock()
136
84
- for _, addr := range p.Addresses {
137
+ for _, addr := range p.addresses {
138
if addr.Equal(a) {
139
return
140
}
141
}
89
- p.Addresses = append(p.Addresses, a)
142
+ p.addresses = append(p.addresses, a)
143
}
144
145
// NetAddress returns the first Multiaddr found for a given network.
93
-func (p *Peer) NetAddress(n string) ma.Multiaddr {
146
+func (p *peer) NetAddress(n string) ma.Multiaddr {
147
p.RLock()
148
defer p.RUnlock()
149
97
- for _, a := range p.Addresses {
150
+ for _, a := range p.addresses {
151
for _, p := range a.Protocols() {
152
if p.Name == n {
153
return a
@@ -105,7 +158,7 @@ func (p *Peer) NetAddress(n string) ma.Multiaddr {
158
}
159
160
// GetLatency retrieves the current latency measurement.
108
-func (p *Peer) GetLatency() (out time.Duration) {
161
+func (p *peer) GetLatency() (out time.Duration) {
162
p.RLock()
163
out = p.latency
164
p.RUnlock()
@@ -116,7 +169,7 @@ func (p *Peer) GetLatency() (out time.Duration) {
169
// TODO: Instead of just keeping a single number,
170
// keep a running average over the last hour or so
171
// Yep, should be EWMA or something. (-jbenet)
119
-func (p *Peer) SetLatency(laten time.Duration) {
172
+func (p *peer) SetLatency(laten time.Duration) {
173
p.Lock()
174
if p.latency == 0 {
175
p.latency = laten
@@ -128,61 +181,98 @@ func (p *Peer) SetLatency(laten time.Duration) {
181
182
// LoadAndVerifyKeyPair unmarshalls, loads a private/public key pair.
183
// Error if (a) unmarshalling fails, or (b) pubkey does not match id.
131
-func (p *Peer) LoadAndVerifyKeyPair(marshalled []byte) error {
184
+func (p *peer) LoadAndVerifyKeyPair(marshalled []byte) error {
185
186
sk, err := ic.UnmarshalPrivateKey(marshalled)
187
if err != nil {
188
return fmt.Errorf("Failed to unmarshal private key: %v", err)
189
}
190
191
+ return p.VerifyAndSetPrivKey(sk)
192
+}
193
+
194
+// VerifyAndSetPrivKey sets private key, given its pubkey matches the peer.ID
195
+func (p *peer) VerifyAndSetPrivKey(sk ic.PrivKey) error {
196
+
197
// construct and assign pubkey. ensure it matches this peer
198
if err := p.VerifyAndSetPubKey(sk.GetPublic()); err != nil {
199
return err
200
}
201
202
// if we didn't have the priavte key, assign it
144
- if p.PrivKey == nil {
145
- p.PrivKey = sk
203
+ if p.privKey == nil {
204
+ p.privKey = sk
205
return nil
206
}
207
208
// if we already had the keys, check they're equal.
150
- if p.PrivKey.Equals(sk) {
209
+ if p.privKey.Equals(sk) {
210
return nil // as expected. keep the old objects.
211
}
212
213
// keys not equal. invariant violated. this warrants a panic.
214
// these keys should be _the same_ because peer.ID = H(pk)
215
// this mismatch should never happen.
157
- log.Error("%s had PrivKey: %v -- got %v", p, p.PrivKey, sk)
216
+ log.Error("%s had PrivKey: %v -- got %v", p, p.privKey, sk)
217
panic("invariant violated: unexpected key mismatch")
218
}
219
220
// VerifyAndSetPubKey sets public key, given it matches the peer.ID
162
-func (p *Peer) VerifyAndSetPubKey(pk ic.PubKey) error {
221
+func (p *peer) VerifyAndSetPubKey(pk ic.PubKey) error {
222
pkid, err := IDFromPubKey(pk)
223
if err != nil {
224
return fmt.Errorf("Failed to hash public key: %v", err)
225
}
226
168
- if !p.ID.Equal(pkid) {
227
+ if !p.id.Equal(pkid) {
228
return fmt.Errorf("Public key does not match peer.ID.")
229
}
230
231
// if we didn't have the keys, assign them.
173
- if p.PubKey == nil {
174
- p.PubKey = pk
232
+ if p.pubKey == nil {
233
+ p.pubKey = pk
234
return nil
235
}
236
237
// if we already had the pubkey, check they're equal.
179
- if p.PubKey.Equals(pk) {
238
+ if p.pubKey.Equals(pk) {
239
return nil // as expected. keep the old objects.
240
}
241
242
// keys not equal. invariant violated. this warrants a panic.
243
// these keys should be _the same_ because peer.ID = H(pk)
244
// this mismatch should never happen.
186
- log.Error("%s had PubKey: %v -- got %v", p, p.PubKey, pk)
245
+ log.Error("%s had PubKey: %v -- got %v", p, p.pubKey, pk)
246
panic("invariant violated: unexpected key mismatch")
247
}
248
+
249
+// WithKeyPair returns a Peer object with given keys.
250
+func WithKeyPair(sk ic.PrivKey, pk ic.PubKey) (Peer, error) {
251
+ if sk == nil && pk == nil {
252
+ return nil, fmt.Errorf("PeerWithKeyPair nil keys")
253
+ }
254
+
255
+ pk2 := sk.GetPublic()
256
+ if pk == nil {
257
+ pk = pk2
258
+ } else if !pk.Equals(pk2) {
259
+ return nil, fmt.Errorf("key mismatch. pubkey is not privkey's pubkey")
260
+ }
261
+
262
+ pkid, err := IDFromPubKey(pk)
263
+ if err != nil {
264
+ return nil, fmt.Errorf("Failed to hash public key: %v", err)
265
+ }
266
+
267
+ return &peer{id: pkid, pubKey: pk, privKey: sk}, nil
268
+}
269
+
270
+// WithID constructs a peer with given ID.
271
+func WithID(id ID) Peer {
272
+ return &peer{id: id}
273
+}
274
+
275
+// WithIDString constructs a peer with given ID (string).
276
+func WithIDString(id string) Peer {
277
+ return WithID(ID(id))
278
+}
peer/peer_test.go
+2
-2
@@ -27,12 +27,12 @@ func TestNetAddress(t *testing.T) {
27
return
28
}
29
30
- p := Peer{ID: ID(mh)}
30
+ p := WithID(ID(mh))
31
p.AddAddress(tcp)
32
p.AddAddress(udp)
33
p.AddAddress(tcp)
34
35
- if len(p.Addresses) == 3 {
35
+ if len(p.Addresses()) == 3 {
36
t.Error("added same address twice")
37
}
38
peer/peerstore.go
+15
-11
@@ -11,8 +11,8 @@ import (
11
12
// Peerstore provides a threadsafe collection for peers.
13
type Peerstore interface {
14
- Get(ID) (*Peer, error)
15
- Put(*Peer) error
14
+ Get(ID) (Peer, error)
15
+ Put(Peer) error
16
Delete(ID) error
17
All() (*Map, error)
18
}
@@ -29,9 +29,13 @@ func NewPeerstore() Peerstore {
29
}
30
}
31
32
-func (p *peerstore) Get(i ID) (*Peer, error) {
33
- p.RLock()
34
- defer p.RUnlock()
32
+func (p *peerstore) Get(i ID) (Peer, error) {
33
+ p.Lock()
34
+ defer p.Unlock()
35
+
36
+ if i == nil {
37
+ panic("wat")
38
+ }
39
40
k := u.Key(i).DsKey()
41
val, err := p.peers.Get(k)
@@ -43,7 +47,7 @@ func (p *peerstore) Get(i ID) (*Peer, error) {
47
48
// not found, construct it ourselves, add it to datastore, and return.
49
case ds.ErrNotFound:
46
- peer := &Peer{ID: i}
50
+ peer := &peer{id: i}
51
if err := p.peers.Put(k, peer); err != nil {
52
return nil, err
53
}
@@ -51,7 +55,7 @@ func (p *peerstore) Get(i ID) (*Peer, error) {
55
56
// no error, got it back fine
57
case nil:
54
- peer, ok := val.(*Peer)
58
+ peer, ok := val.(*peer)
59
if !ok {
60
return nil, errors.New("stored value was not a Peer")
61
}
@@ -59,11 +63,11 @@ func (p *peerstore) Get(i ID) (*Peer, error) {
63
}
64
}
65
62
-func (p *peerstore) Put(peer *Peer) error {
66
+func (p *peerstore) Put(peer Peer) error {
67
p.Lock()
68
defer p.Unlock()
69
66
- k := u.Key(peer.ID).DsKey()
70
+ k := peer.Key().DsKey()
71
return p.peers.Put(k, peer)
72
}
73
@@ -91,9 +95,9 @@ func (p *peerstore) All() (*Map, error) {
95
continue
96
}
97
94
- pval, ok := val.(*Peer)
98
+ pval, ok := val.(*peer)
99
if ok {
96
- (*ps)[u.Key(pval.ID)] = pval
100
+ (*ps)[pval.Key()] = pval
101
}
102
}
103
return ps, nil
peer/peerstore_test.go
+2
-2
@@ -7,13 +7,13 @@ import (
7
ma "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-multiaddr"
8
)
9
10
-func setupPeer(id string, addr string) (*Peer, error) {
10
+func setupPeer(id string, addr string) (Peer, error) {
11
tcp, err := ma.NewMultiaddr(addr)
12
if err != nil {
13
return nil, err
14
}
15
16
- p := &Peer{ID: ID(id)}
16
+ p := WithIDString(id)
17
p.AddAddress(tcp)
18
return p, nil
19
}
peer/queue/distance.go
+4
-4
@@ -13,7 +13,7 @@ import (
13
// peerMetric tracks a peer and its distance to something else.
14
type peerMetric struct {
15
// the peer
16
- peer *peer.Peer
16
+ peer peer.Peer
17
18
// big.Int for XOR metric
19
metric *big.Int
@@ -64,11 +64,11 @@ func (pq *distancePQ) Len() int {
64
return len(pq.heap)
65
}
66
67
-func (pq *distancePQ) Enqueue(p *peer.Peer) {
67
+func (pq *distancePQ) Enqueue(p peer.Peer) {
68
pq.Lock()
69
defer pq.Unlock()
70
71
- distance := ks.XORKeySpace.Key(p.ID).Distance(pq.from)
71
+ distance := ks.XORKeySpace.Key(p.ID()).Distance(pq.from)
72
73
heap.Push(&pq.heap, &peerMetric{
74
peer: p,
@@ -76,7 +76,7 @@ func (pq *distancePQ) Enqueue(p *peer.Peer) {
76
})
77
}
78
79
-func (pq *distancePQ) Dequeue() *peer.Peer {
79
+func (pq *distancePQ) Dequeue() peer.Peer {
80
pq.Lock()
81
defer pq.Unlock()
82
peer/queue/interface.go
+2
-2
@@ -11,8 +11,8 @@ type PeerQueue interface {
11
Len() int
12
13
// Enqueue adds this node to the queue.
14
- Enqueue(*peer.Peer)
14
+ Enqueue(peer.Peer)
15
16
// Dequeue retrieves the highest (smallest int) priority node
17
- Dequeue() *peer.Peer
17
+ Dequeue() peer.Peer
18
}
peer/queue/queue_test.go
+4
-4
@@ -12,8 +12,8 @@ import (
12
context "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
13
)
14
15
-func newPeer(id string) *peer.Peer {
16
- return &peer.Peer{ID: peer.ID(id)}
15
+func newPeer(id string) peer.Peer {
16
+ return peer.WithIDString(id)
17
}
18
19
func TestQueue(t *testing.T) {
@@ -66,10 +66,10 @@ func TestQueue(t *testing.T) {
66
67
}
68
69
-func newPeerTime(t time.Time) *peer.Peer {
69
+func newPeerTime(t time.Time) peer.Peer {
70
s := fmt.Sprintf("hmmm time: %v", t)
71
h := u.Hash([]byte(s))
72
- return &peer.Peer{ID: peer.ID(h)}
72
+ return peer.WithID(peer.ID(h))
73
}
74
75
func TestSyncQueue(t *testing.T) {
peer/queue/sync.go
+6
-6
@@ -9,8 +9,8 @@ import (
9
// ChanQueue makes any PeerQueue synchronizable through channels.
10
type ChanQueue struct {
11
Queue PeerQueue
12
- EnqChan chan<- *peer.Peer
13
- DeqChan <-chan *peer.Peer
12
+ EnqChan chan<- peer.Peer
13
+ DeqChan <-chan peer.Peer
14
}
15
16
// NewChanQueue creates a ChanQueue by wrapping pq.
@@ -23,8 +23,8 @@ func NewChanQueue(ctx context.Context, pq PeerQueue) *ChanQueue {
23
func (cq *ChanQueue) process(ctx context.Context) {
24
25
// construct the channels here to be able to use them bidirectionally
26
- enqChan := make(chan *peer.Peer, 10)
27
- deqChan := make(chan *peer.Peer, 10)
26
+ enqChan := make(chan peer.Peer, 10)
27
+ deqChan := make(chan peer.Peer, 10)
28
29
cq.EnqChan = enqChan
30
cq.DeqChan = deqChan
@@ -32,8 +32,8 @@ func (cq *ChanQueue) process(ctx context.Context) {
32
go func() {
33
defer close(deqChan)
34
35
- var next *peer.Peer
36
- var item *peer.Peer
35
+ var next peer.Peer
36
+ var item peer.Peer
37
var more bool
38
39
for {
routing/dht/Message.go
+6
-5
@@ -17,20 +17,21 @@ func newMessage(typ Message_MessageType, key string, level int) *Message {
17
return m
18
}
19
20
-func peerToPBPeer(p *peer.Peer) *Message_Peer {
20
+func peerToPBPeer(p peer.Peer) *Message_Peer {
21
pbp := new(Message_Peer)
22
- if len(p.Addresses) == 0 || p.Addresses[0] == nil {
22
+ addrs := p.Addresses()
23
+ if len(addrs) == 0 || addrs[0] == nil {
24
pbp.Addr = proto.String("")
25
} else {
25
- addr := p.Addresses[0].String()
26
+ addr := addrs[0].String()
27
pbp.Addr = &addr
28
}
28
- pid := string(p.ID)
29
+ pid := string(p.ID())
30
pbp.Id = &pid
31
return pbp
32
}
33
33
-func peersToPBPeers(peers []*peer.Peer) []*Message_Peer {
34
+func peersToPBPeers(peers []peer.Peer) []*Message_Peer {
35
pbpeers := make([]*Message_Peer, len(peers))
36
for i, p := range peers {
37
pbpeers[i] = peerToPBPeer(p)
routing/dht/dht.go
+33
-33
@@ -39,7 +39,7 @@ type IpfsDHT struct {
39
sender inet.Sender
40
41
// Local peer (yourself)
42
- self *peer.Peer
42
+ self peer.Peer
43
44
// Other peers
45
peerstore peer.Peerstore
@@ -60,7 +60,7 @@ type IpfsDHT struct {
60
}
61
62
// NewDHT creates a new DHT object with the given peer as the 'local' host
63
-func NewDHT(ctx context.Context, p *peer.Peer, ps peer.Peerstore, net inet.Network, sender inet.Sender, dstore ds.Datastore) *IpfsDHT {
63
+func NewDHT(ctx context.Context, p peer.Peer, ps peer.Peerstore, net inet.Network, sender inet.Sender, dstore ds.Datastore) *IpfsDHT {
64
dht := new(IpfsDHT)
65
dht.network = net
66
dht.sender = sender
@@ -69,12 +69,12 @@ func NewDHT(ctx context.Context, p *peer.Peer, ps peer.Peerstore, net inet.Netwo
69
dht.peerstore = ps
70
dht.ctx = ctx
71
72
- dht.providers = NewProviderManager(p.ID)
72
+ dht.providers = NewProviderManager(p.ID())
73
74
dht.routingTables = make([]*kb.RoutingTable, 3)
75
- dht.routingTables[0] = kb.NewRoutingTable(20, kb.ConvertPeerID(p.ID), time.Millisecond*1000)
76
- dht.routingTables[1] = kb.NewRoutingTable(20, kb.ConvertPeerID(p.ID), time.Millisecond*1000)
77
- dht.routingTables[2] = kb.NewRoutingTable(20, kb.ConvertPeerID(p.ID), time.Hour)
75
+ dht.routingTables[0] = kb.NewRoutingTable(20, kb.ConvertPeerID(p.ID()), time.Millisecond*1000)
76
+ dht.routingTables[1] = kb.NewRoutingTable(20, kb.ConvertPeerID(p.ID()), time.Millisecond*1000)
77
+ dht.routingTables[2] = kb.NewRoutingTable(20, kb.ConvertPeerID(p.ID()), time.Hour)
78
dht.birth = time.Now()
79
80
if doPinging {
@@ -84,7 +84,7 @@ func NewDHT(ctx context.Context, p *peer.Peer, ps peer.Peerstore, net inet.Netwo
84
}
85
86
// Connect to a new peer at the given address, ping and add to the routing table
87
-func (dht *IpfsDHT) Connect(ctx context.Context, npeer *peer.Peer) (*peer.Peer, error) {
87
+func (dht *IpfsDHT) Connect(ctx context.Context, npeer peer.Peer) (peer.Peer, error) {
88
log.Debug("Connect to new peer: %s", npeer)
89
90
// TODO(jbenet,whyrusleeping)
@@ -175,7 +175,7 @@ func (dht *IpfsDHT) HandleMessage(ctx context.Context, mes msg.NetMessage) msg.N
175
176
// sendRequest sends out a request using dht.sender, but also makes sure to
177
// measure the RTT for latency measurements.
178
-func (dht *IpfsDHT) sendRequest(ctx context.Context, p *peer.Peer, pmes *Message) (*Message, error) {
178
+func (dht *IpfsDHT) sendRequest(ctx context.Context, p peer.Peer, pmes *Message) (*Message, error) {
179
180
mes, err := msg.FromObject(p, pmes)
181
if err != nil {
@@ -208,7 +208,7 @@ func (dht *IpfsDHT) sendRequest(ctx context.Context, p *peer.Peer, pmes *Message
208
}
209
210
// putValueToNetwork stores the given key/value pair at the peer 'p'
211
-func (dht *IpfsDHT) putValueToNetwork(ctx context.Context, p *peer.Peer,
211
+func (dht *IpfsDHT) putValueToNetwork(ctx context.Context, p peer.Peer,
212
key string, value []byte) error {
213
214
pmes := newMessage(Message_PUT_VALUE, string(key), 0)
@@ -224,12 +224,12 @@ func (dht *IpfsDHT) putValueToNetwork(ctx context.Context, p *peer.Peer,
224
return nil
225
}
226
227
-func (dht *IpfsDHT) putProvider(ctx context.Context, p *peer.Peer, key string) error {
227
+func (dht *IpfsDHT) putProvider(ctx context.Context, p peer.Peer, key string) error {
228
229
pmes := newMessage(Message_ADD_PROVIDER, string(key), 0)
230
231
// add self as the provider
232
- pmes.ProviderPeers = peersToPBPeers([]*peer.Peer{dht.self})
232
+ pmes.ProviderPeers = peersToPBPeers([]peer.Peer{dht.self})
233
234
rpmes, err := dht.sendRequest(ctx, p, pmes)
235
if err != nil {
@@ -244,8 +244,8 @@ func (dht *IpfsDHT) putProvider(ctx context.Context, p *peer.Peer, key string) e
244
return nil
245
}
246
247
-func (dht *IpfsDHT) getValueOrPeers(ctx context.Context, p *peer.Peer,
248
- key u.Key, level int) ([]byte, []*peer.Peer, error) {
247
+func (dht *IpfsDHT) getValueOrPeers(ctx context.Context, p peer.Peer,
248
+ key u.Key, level int) ([]byte, []peer.Peer, error) {
249
250
pmes, err := dht.getValueSingle(ctx, p, key, level)
251
if err != nil {
@@ -270,7 +270,7 @@ func (dht *IpfsDHT) getValueOrPeers(ctx context.Context, p *peer.Peer,
270
}
271
272
// Perhaps we were given closer peers
273
- var peers []*peer.Peer
273
+ var peers []peer.Peer
274
for _, pb := range pmes.GetCloserPeers() {
275
pr, err := dht.addPeer(pb)
276
if err != nil {
@@ -289,8 +289,8 @@ func (dht *IpfsDHT) getValueOrPeers(ctx context.Context, p *peer.Peer,
289
return nil, nil, u.ErrNotFound
290
}
291
292
-func (dht *IpfsDHT) addPeer(pb *Message_Peer) (*peer.Peer, error) {
293
- if peer.ID(pb.GetId()).Equal(dht.self.ID) {
292
+func (dht *IpfsDHT) addPeer(pb *Message_Peer) (peer.Peer, error) {
293
+ if peer.ID(pb.GetId()).Equal(dht.self.ID()) {
294
return nil, errors.New("cannot add self as peer")
295
}
296
@@ -310,7 +310,7 @@ func (dht *IpfsDHT) addPeer(pb *Message_Peer) (*peer.Peer, error) {
310
}
311
312
// getValueSingle simply performs the get value RPC with the given parameters
313
-func (dht *IpfsDHT) getValueSingle(ctx context.Context, p *peer.Peer,
313
+func (dht *IpfsDHT) getValueSingle(ctx context.Context, p peer.Peer,
314
key u.Key, level int) (*Message, error) {
315
316
pmes := newMessage(Message_GET_VALUE, string(key), level)
@@ -369,7 +369,7 @@ func (dht *IpfsDHT) putLocal(key u.Key, value []byte) error {
369
370
// Update signals to all routingTables to Update their last-seen status
371
// on the given peer.
372
-func (dht *IpfsDHT) Update(p *peer.Peer) {
372
+func (dht *IpfsDHT) Update(p peer.Peer) {
373
log.Debug("updating peer: %s latency = %f\n", p, p.GetLatency().Seconds())
374
removedCount := 0
375
for _, route := range dht.routingTables {
@@ -390,7 +390,7 @@ func (dht *IpfsDHT) Update(p *peer.Peer) {
390
}
391
392
// FindLocal looks for a peer with a given ID connected to this dht and returns the peer and the table it was found in.
393
-func (dht *IpfsDHT) FindLocal(id peer.ID) (*peer.Peer, *kb.RoutingTable) {
393
+func (dht *IpfsDHT) FindLocal(id peer.ID) (peer.Peer, *kb.RoutingTable) {
394
for _, table := range dht.routingTables {
395
p := table.Find(id)
396
if p != nil {
@@ -400,7 +400,7 @@ func (dht *IpfsDHT) FindLocal(id peer.ID) (*peer.Peer, *kb.RoutingTable) {
400
return nil, nil
401
}
402
403
-func (dht *IpfsDHT) findPeerSingle(ctx context.Context, p *peer.Peer, id peer.ID, level int) (*Message, error) {
403
+func (dht *IpfsDHT) findPeerSingle(ctx context.Context, p peer.Peer, id peer.ID, level int) (*Message, error) {
404
pmes := newMessage(Message_FIND_NODE, string(id), level)
405
return dht.sendRequest(ctx, p, pmes)
406
}
@@ -411,14 +411,14 @@ func (dht *IpfsDHT) printTables() {
411
}
412
}
413
414
-func (dht *IpfsDHT) findProvidersSingle(ctx context.Context, p *peer.Peer, key u.Key, level int) (*Message, error) {
414
+func (dht *IpfsDHT) findProvidersSingle(ctx context.Context, p peer.Peer, key u.Key, level int) (*Message, error) {
415
pmes := newMessage(Message_GET_PROVIDERS, string(key), level)
416
return dht.sendRequest(ctx, p, pmes)
417
}
418
419
// TODO: Could be done async
420
-func (dht *IpfsDHT) addProviders(key u.Key, peers []*Message_Peer) []*peer.Peer {
421
- var provArr []*peer.Peer
420
+func (dht *IpfsDHT) addProviders(key u.Key, peers []*Message_Peer) []peer.Peer {
421
+ var provArr []peer.Peer
422
for _, prov := range peers {
423
p, err := dht.peerFromInfo(prov)
424
if err != nil {
@@ -429,7 +429,7 @@ func (dht *IpfsDHT) addProviders(key u.Key, peers []*Message_Peer) []*peer.Peer
429
log.Debug("%s adding provider: %s for %s", dht.self, p, key)
430
431
// Dont add outselves to the list
432
- if p.ID.Equal(dht.self.ID) {
432
+ if p.ID().Equal(dht.self.ID()) {
433
continue
434
}
435
@@ -441,7 +441,7 @@ func (dht *IpfsDHT) addProviders(key u.Key, peers []*Message_Peer) []*peer.Peer
441
}
442
443
// nearestPeersToQuery returns the routing tables closest peers.
444
-func (dht *IpfsDHT) nearestPeersToQuery(pmes *Message, count int) []*peer.Peer {
444
+func (dht *IpfsDHT) nearestPeersToQuery(pmes *Message, count int) []peer.Peer {
445
level := pmes.GetClusterLevel()
446
cluster := dht.routingTables[level]
447
@@ -451,7 +451,7 @@ func (dht *IpfsDHT) nearestPeersToQuery(pmes *Message, count int) []*peer.Peer {
451
}
452
453
// betterPeerToQuery returns nearestPeersToQuery, but iff closer than self.
454
-func (dht *IpfsDHT) betterPeersToQuery(pmes *Message, count int) []*peer.Peer {
454
+func (dht *IpfsDHT) betterPeersToQuery(pmes *Message, count int) []peer.Peer {
455
closer := dht.nearestPeersToQuery(pmes, count)
456
457
// no node? nil
@@ -461,17 +461,17 @@ func (dht *IpfsDHT) betterPeersToQuery(pmes *Message, count int) []*peer.Peer {
461
462
// == to self? thats bad
463
for _, p := range closer {
464
- if p.ID.Equal(dht.self.ID) {
464
+ if p.ID().Equal(dht.self.ID()) {
465
log.Error("Attempted to return self! this shouldnt happen...")
466
return nil
467
}
468
}
469
470
- var filtered []*peer.Peer
470
+ var filtered []peer.Peer
471
for _, p := range closer {
472
// must all be closer than self
473
key := u.Key(pmes.GetKey())
474
- if !kb.Closer(dht.self.ID, p.ID, key) {
474
+ if !kb.Closer(dht.self.ID(), p.ID(), key) {
475
filtered = append(filtered, p)
476
}
477
}
@@ -480,7 +480,7 @@ func (dht *IpfsDHT) betterPeersToQuery(pmes *Message, count int) []*peer.Peer {
480
return filtered
481
}
482
483
-func (dht *IpfsDHT) getPeer(id peer.ID) (*peer.Peer, error) {
483
+func (dht *IpfsDHT) getPeer(id peer.ID) (peer.Peer, error) {
484
p, err := dht.peerstore.Get(id)
485
if err != nil {
486
err = fmt.Errorf("Failed to get peer from peerstore: %s", err)
@@ -490,12 +490,12 @@ func (dht *IpfsDHT) getPeer(id peer.ID) (*peer.Peer, error) {
490
return p, nil
491
}
492
493
-func (dht *IpfsDHT) peerFromInfo(pbp *Message_Peer) (*peer.Peer, error) {
493
+func (dht *IpfsDHT) peerFromInfo(pbp *Message_Peer) (peer.Peer, error) {
494
495
id := peer.ID(pbp.GetId())
496
497
// continue if it's ourselves
498
- if id.Equal(dht.self.ID) {
498
+ if id.Equal(dht.self.ID()) {
499
return nil, errors.New("found self")
500
}
501
@@ -512,7 +512,7 @@ func (dht *IpfsDHT) peerFromInfo(pbp *Message_Peer) (*peer.Peer, error) {
512
return p, nil
513
}
514
515
-func (dht *IpfsDHT) ensureConnectedToPeer(pbp *Message_Peer) (*peer.Peer, error) {
515
+func (dht *IpfsDHT) ensureConnectedToPeer(pbp *Message_Peer) (peer.Peer, error) {
516
p, err := dht.peerFromInfo(pbp)
517
if err != nil {
518
return nil, err
routing/dht/dht_test.go
+9
-14
@@ -20,7 +20,7 @@ import (
20
"time"
21
)
22
23
-func setupDHT(ctx context.Context, t *testing.T, p *peer.Peer) *IpfsDHT {
23
+func setupDHT(ctx context.Context, t *testing.T, p peer.Peer) *IpfsDHT {
24
peerstore := peer.NewPeerstore()
25
26
dhts := netservice.NewService(nil) // nil handler for now, need to patch it
@@ -40,7 +40,7 @@ func setupDHT(ctx context.Context, t *testing.T, p *peer.Peer) *IpfsDHT {
40
return d
41
}
42
43
-func setupDHTS(ctx context.Context, n int, t *testing.T) ([]ma.Multiaddr, []*peer.Peer, []*IpfsDHT) {
43
+func setupDHTS(ctx context.Context, n int, t *testing.T) ([]ma.Multiaddr, []peer.Peer, []*IpfsDHT) {
44
var addrs []ma.Multiaddr
45
for i := 0; i < n; i++ {
46
a, err := ma.NewMultiaddr(fmt.Sprintf("/ip4/127.0.0.1/tcp/%d", 5000+i))
@@ -50,7 +50,7 @@ func setupDHTS(ctx context.Context, n int, t *testing.T) ([]ma.Multiaddr, []*pee
50
addrs = append(addrs, a)
51
}
52
53
- var peers []*peer.Peer
53
+ var peers []peer.Peer
54
for i := 0; i < n; i++ {
55
p := makePeer(addrs[i])
56
peers = append(peers, p)
@@ -64,21 +64,16 @@ func setupDHTS(ctx context.Context, n int, t *testing.T) ([]ma.Multiaddr, []*pee
64
return addrs, peers, dhts
65
}
66
67
-func makePeer(addr ma.Multiaddr) *peer.Peer {
68
- p := new(peer.Peer)
69
- p.AddAddress(addr)
67
+func makePeer(addr ma.Multiaddr) peer.Peer {
68
sk, pk, err := ci.GenerateKeyPair(ci.RSA, 512)
69
if err != nil {
70
panic(err)
71
}
74
- p.PrivKey = sk
75
- p.PubKey = pk
76
- id, err := peer.IDFromPubKey(pk)
72
+ p, err := peer.WithKeyPair(sk, pk)
73
if err != nil {
74
panic(err)
75
}
80
-
81
- p.ID = id
76
+ p.AddAddress(addr)
77
return p
78
}
79
@@ -289,7 +284,7 @@ func TestProvidesAsync(t *testing.T) {
284
provs := dhts[0].FindProvidersAsync(ctxT, u.Key("hello"), 5)
285
select {
286
case p := <-provs:
292
- if !p.ID.Equal(dhts[3].self.ID) {
287
+ if !p.ID().Equal(dhts[3].self.ID()) {
288
t.Fatalf("got a provider, but not the right one. %s", p)
289
}
290
case <-ctxT.Done():
@@ -379,7 +374,7 @@ func TestFindPeer(t *testing.T) {
374
}
375
376
ctxT, _ := context.WithTimeout(ctx, time.Second)
382
- p, err := dhts[0].FindPeer(ctxT, peers[2].ID)
377
+ p, err := dhts[0].FindPeer(ctxT, peers[2].ID())
378
if err != nil {
379
t.Fatal(err)
380
}
@@ -388,7 +383,7 @@ func TestFindPeer(t *testing.T) {
383
t.Fatal("Failed to find peer.")
384
}
385
391
- if !p.ID.Equal(peers[2].ID) {
386
+ if !p.ID().Equal(peers[2].ID()) {
387
t.Fatal("Didnt find expected peer.")
388
}
389
}
routing/dht/diag.go
+3
-2
@@ -32,12 +32,13 @@ func (di *diagInfo) Marshal() []byte {
32
func (dht *IpfsDHT) getDiagInfo() *diagInfo {
33
di := new(diagInfo)
34
di.CodeVersion = "github.com/jbenet/go-ipfs"
35
- di.ID = dht.self.ID
35
+ di.ID = dht.self.ID()
36
di.LifeSpan = time.Since(dht.birth)
37
di.Keys = nil // Currently no way to query datastore
38
39
for _, p := range dht.routingTables[0].ListPeers() {
40
- di.Connections = append(di.Connections, connDiagInfo{p.GetLatency(), p.ID})
40
+ d := connDiagInfo{p.GetLatency(), p.ID()}
41
+ di.Connections = append(di.Connections, d)
42
}
43
return di
44
}
routing/dht/ext_test.go
+18
-21
@@ -9,7 +9,6 @@ import (
9
"github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/goprotobuf/proto"
10
11
ds "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/datastore.go"
12
- ma "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-multiaddr"
12
msg "github.com/jbenet/go-ipfs/net/message"
13
mux "github.com/jbenet/go-ipfs/net/mux"
14
peer "github.com/jbenet/go-ipfs/peer"
@@ -66,17 +65,17 @@ type fauxNet struct {
65
}
66
67
// DialPeer attempts to establish a connection to a given peer
69
-func (f *fauxNet) DialPeer(*peer.Peer) error {
68
+func (f *fauxNet) DialPeer(peer.Peer) error {
69
return nil
70
}
71
72
// ClosePeer connection to peer
74
-func (f *fauxNet) ClosePeer(*peer.Peer) error {
73
+func (f *fauxNet) ClosePeer(peer.Peer) error {
74
return nil
75
}
76
77
// IsConnected returns whether a connection to given peer exists.
79
-func (f *fauxNet) IsConnected(*peer.Peer) (bool, error) {
78
+func (f *fauxNet) IsConnected(peer.Peer) (bool, error) {
79
return true, nil
80
}
81
@@ -88,7 +87,7 @@ func (f *fauxNet) SendMessage(msg.NetMessage) error {
87
return nil
88
}
89
91
-func (f *fauxNet) GetPeerList() []*peer.Peer {
90
+func (f *fauxNet) GetPeerList() []peer.Peer {
91
return nil
92
}
93
@@ -107,11 +106,10 @@ func TestGetFailures(t *testing.T) {
106
fs := &fauxSender{}
107
108
peerstore := peer.NewPeerstore()
110
- local := new(peer.Peer)
111
- local.ID = peer.ID("test_peer")
109
+ local := peer.WithIDString("test_peer")
110
111
d := NewDHT(ctx, local, peerstore, fn, fs, ds.NewMapDatastore())
114
- other := &peer.Peer{ID: peer.ID("other_peer")}
112
+ other := peer.WithIDString("other_peer")
113
d.Update(other)
114
115
// This one should time out
@@ -189,11 +187,10 @@ func TestGetFailures(t *testing.T) {
187
}
188
189
// TODO: Maybe put these in some sort of "ipfs_testutil" package
192
-func _randPeer() *peer.Peer {
193
- p := new(peer.Peer)
194
- p.ID = make(peer.ID, 16)
195
- p.Addresses = []ma.Multiaddr{nil}
196
- crand.Read(p.ID)
190
+func _randPeer() peer.Peer {
191
+ id := make(peer.ID, 16)
192
+ crand.Read(id)
193
+ p := peer.WithID(id)
194
return p
195
}
196
@@ -204,13 +201,13 @@ func TestNotFound(t *testing.T) {
201
fn := &fauxNet{}
202
fs := &fauxSender{}
203
207
- local := new(peer.Peer)
208
- local.ID = peer.ID("test_peer")
204
+ local := peer.WithIDString("test_peer")
205
peerstore := peer.NewPeerstore()
206
+ peerstore.Put(local)
207
208
d := NewDHT(ctx, local, peerstore, fn, fs, ds.NewMapDatastore())
209
213
- var ps []*peer.Peer
210
+ var ps []peer.Peer
211
for i := 0; i < 5; i++ {
212
ps = append(ps, _randPeer())
213
d.Update(ps[i])
@@ -228,7 +225,7 @@ func TestNotFound(t *testing.T) {
225
case Message_GET_VALUE:
226
resp := &Message{Type: pmes.Type}
227
231
- peers := []*peer.Peer{}
228
+ peers := []peer.Peer{}
229
for i := 0; i < 7; i++ {
230
peers = append(peers, _randPeer())
231
}
@@ -270,13 +267,13 @@ func TestLessThanKResponses(t *testing.T) {
267
u.Debug = false
268
fn := &fauxNet{}
269
fs := &fauxSender{}
270
+ local := peer.WithIDString("test_peer")
271
peerstore := peer.NewPeerstore()
274
- local := new(peer.Peer)
275
- local.ID = peer.ID("test_peer")
272
+ peerstore.Put(local)
273
274
d := NewDHT(ctx, local, peerstore, fn, fs, ds.NewMapDatastore())
275
279
- var ps []*peer.Peer
276
+ var ps []peer.Peer
277
for i := 0; i < 5; i++ {
278
ps = append(ps, _randPeer())
279
d.Update(ps[i])
@@ -295,7 +292,7 @@ func TestLessThanKResponses(t *testing.T) {
292
case Message_GET_VALUE:
293
resp := &Message{
294
Type: pmes.Type,
298
- CloserPeers: peersToPBPeers([]*peer.Peer{other}),
295
+ CloserPeers: peersToPBPeers([]peer.Peer{other}),
296
}
297
298
mes, err := msg.FromObject(mes.Peer(), resp)
routing/dht/handlers.go
+14
-14
@@ -14,7 +14,7 @@ import (
14
var CloserPeerCount = 4
15
16
// dhthandler specifies the signature of functions that handle DHT messages.
17
-type dhtHandler func(*peer.Peer, *Message) (*Message, error)
17
+type dhtHandler func(peer.Peer, *Message) (*Message, error)
18
19
func (dht *IpfsDHT) handlerForMsgType(t Message_MessageType) dhtHandler {
20
switch t {
@@ -35,7 +35,7 @@ func (dht *IpfsDHT) handlerForMsgType(t Message_MessageType) dhtHandler {
35
}
36
}
37
38
-func (dht *IpfsDHT) handleGetValue(p *peer.Peer, pmes *Message) (*Message, error) {
38
+func (dht *IpfsDHT) handleGetValue(p peer.Peer, pmes *Message) (*Message, error) {
39
log.Debug("%s handleGetValue for key: %s\n", dht.self, pmes.GetKey())
40
41
// setup response
@@ -93,7 +93,7 @@ func (dht *IpfsDHT) handleGetValue(p *peer.Peer, pmes *Message) (*Message, error
93
}
94
95
// Store a value in this peer local storage
96
-func (dht *IpfsDHT) handlePutValue(p *peer.Peer, pmes *Message) (*Message, error) {
96
+func (dht *IpfsDHT) handlePutValue(p peer.Peer, pmes *Message) (*Message, error) {
97
dht.dslock.Lock()
98
defer dht.dslock.Unlock()
99
dskey := u.Key(pmes.GetKey()).DsKey()
@@ -102,18 +102,18 @@ func (dht *IpfsDHT) handlePutValue(p *peer.Peer, pmes *Message) (*Message, error
102
return pmes, err
103
}
104
105
-func (dht *IpfsDHT) handlePing(p *peer.Peer, pmes *Message) (*Message, error) {
105
+func (dht *IpfsDHT) handlePing(p peer.Peer, pmes *Message) (*Message, error) {
106
log.Debug("%s Responding to ping from %s!\n", dht.self, p)
107
return pmes, nil
108
}
109
110
-func (dht *IpfsDHT) handleFindPeer(p *peer.Peer, pmes *Message) (*Message, error) {
110
+func (dht *IpfsDHT) handleFindPeer(p peer.Peer, pmes *Message) (*Message, error) {
111
resp := newMessage(pmes.GetType(), "", pmes.GetClusterLevel())
112
- var closest []*peer.Peer
112
+ var closest []peer.Peer
113
114
// if looking for self... special case where we send it on CloserPeers.
115
- if peer.ID(pmes.GetKey()).Equal(dht.self.ID) {
116
- closest = []*peer.Peer{dht.self}
115
+ if peer.ID(pmes.GetKey()).Equal(dht.self.ID()) {
116
+ closest = []peer.Peer{dht.self}
117
} else {
118
closest = dht.betterPeersToQuery(pmes, CloserPeerCount)
119
}
@@ -123,9 +123,9 @@ func (dht *IpfsDHT) handleFindPeer(p *peer.Peer, pmes *Message) (*Message, error
123
return resp, nil
124
}
125
126
- var withAddresses []*peer.Peer
126
+ var withAddresses []peer.Peer
127
for _, p := range closest {
128
- if len(p.Addresses) > 0 {
128
+ if len(p.Addresses()) > 0 {
129
withAddresses = append(withAddresses, p)
130
}
131
}
@@ -137,7 +137,7 @@ func (dht *IpfsDHT) handleFindPeer(p *peer.Peer, pmes *Message) (*Message, error
137
return resp, nil
138
}
139
140
-func (dht *IpfsDHT) handleGetProviders(p *peer.Peer, pmes *Message) (*Message, error) {
140
+func (dht *IpfsDHT) handleGetProviders(p peer.Peer, pmes *Message) (*Message, error) {
141
resp := newMessage(pmes.GetType(), pmes.GetKey(), pmes.GetClusterLevel())
142
143
// check if we have this value, to add ourselves as provider.
@@ -171,10 +171,10 @@ func (dht *IpfsDHT) handleGetProviders(p *peer.Peer, pmes *Message) (*Message, e
171
172
type providerInfo struct {
173
Creation time.Time
174
- Value *peer.Peer
174
+ Value peer.Peer
175
}
176
177
-func (dht *IpfsDHT) handleAddProvider(p *peer.Peer, pmes *Message) (*Message, error) {
177
+func (dht *IpfsDHT) handleAddProvider(p peer.Peer, pmes *Message) (*Message, error) {
178
key := u.Key(pmes.GetKey())
179
180
log.Debug("%s adding %s as a provider for '%s'\n", dht.self, p, peer.ID(key))
@@ -182,7 +182,7 @@ func (dht *IpfsDHT) handleAddProvider(p *peer.Peer, pmes *Message) (*Message, er
182
// add provider should use the address given in the message
183
for _, pb := range pmes.GetProviderPeers() {
184
pid := peer.ID(pb.GetId())
185
- if pid.Equal(p.ID) {
185
+ if pid.Equal(p.ID()) {
186
187
addr, err := pb.Address()
188
if err != nil {
routing/dht/providers.go
+7
-7
@@ -20,12 +20,12 @@ type ProviderManager struct {
20
21
type addProv struct {
22
k u.Key
23
- val *peer.Peer
23
+ val peer.Peer
24
}
25
26
type getProv struct {
27
k u.Key
28
- resp chan []*peer.Peer
28
+ resp chan []peer.Peer
29
}
30
31
func NewProviderManager(local peer.ID) *ProviderManager {
@@ -45,7 +45,7 @@ func (pm *ProviderManager) run() {
45
for {
46
select {
47
case np := <-pm.newprovs:
48
- if np.val.ID.Equal(pm.lpeer) {
48
+ if np.val.ID().Equal(pm.lpeer) {
49
pm.local[np.k] = struct{}{}
50
}
51
pi := new(providerInfo)
@@ -54,7 +54,7 @@ func (pm *ProviderManager) run() {
54
arr := pm.providers[np.k]
55
pm.providers[np.k] = append(arr, pi)
56
case gp := <-pm.getprovs:
57
- var parr []*peer.Peer
57
+ var parr []peer.Peer
58
provs := pm.providers[gp.k]
59
for _, p := range provs {
60
parr = append(parr, p.Value)
@@ -82,17 +82,17 @@ func (pm *ProviderManager) run() {
82
}
83
}
84
85
-func (pm *ProviderManager) AddProvider(k u.Key, val *peer.Peer) {
85
+func (pm *ProviderManager) AddProvider(k u.Key, val peer.Peer) {
86
pm.newprovs <- &addProv{
87
k: k,
88
val: val,
89
}
90
}
91
92
-func (pm *ProviderManager) GetProviders(k u.Key) []*peer.Peer {
92
+func (pm *ProviderManager) GetProviders(k u.Key) []peer.Peer {
93
gp := new(getProv)
94
gp.k = k
95
- gp.resp = make(chan []*peer.Peer)
95
+ gp.resp = make(chan []peer.Peer)
96
pm.getprovs <- gp
97
return <-gp.resp
98
}
routing/dht/providers_test.go
+1
-1
@@ -11,7 +11,7 @@ func TestProviderManager(t *testing.T) {
11
mid := peer.ID("testing")
12
p := NewProviderManager(mid)
13
a := u.Key("test")
14
- p.AddProvider(a, &peer.Peer{})
14
+ p.AddProvider(a, peer.WithIDString("testingprovider"))
15
resp := p.GetProviders(a)
16
if len(resp) != 1 {
17
t.Fatal("Could not retrieve provider.")
routing/dht/query.go
+10
-10
@@ -26,10 +26,10 @@ type dhtQuery struct {
26
}
27
28
type dhtQueryResult struct {
29
- value []byte // GetValue
30
- peer *peer.Peer // FindPeer
31
- providerPeers []*peer.Peer // GetProviders
32
- closerPeers []*peer.Peer // *
29
+ value []byte // GetValue
30
+ peer peer.Peer // FindPeer
31
+ providerPeers []peer.Peer // GetProviders
32
+ closerPeers []peer.Peer // *
33
success bool
34
}
35
@@ -47,10 +47,10 @@ func newQuery(k u.Key, f queryFunc) *dhtQuery {
47
// - the value
48
// - a list of peers potentially better able to serve the query
49
// - an error
50
-type queryFunc func(context.Context, *peer.Peer) (*dhtQueryResult, error)
50
+type queryFunc func(context.Context, peer.Peer) (*dhtQueryResult, error)
51
52
// Run runs the query at hand. pass in a list of peers to use first.
53
-func (q *dhtQuery) Run(ctx context.Context, peers []*peer.Peer) (*dhtQueryResult, error) {
53
+func (q *dhtQuery) Run(ctx context.Context, peers []peer.Peer) (*dhtQueryResult, error) {
54
runner := newQueryRunner(ctx, q)
55
return runner.Run(peers)
56
}
@@ -100,7 +100,7 @@ func newQueryRunner(ctx context.Context, q *dhtQuery) *dhtQueryRunner {
100
}
101
}
102
103
-func (r *dhtQueryRunner) Run(peers []*peer.Peer) (*dhtQueryResult, error) {
103
+func (r *dhtQueryRunner) Run(peers []peer.Peer) (*dhtQueryResult, error) {
104
log.Debug("Run query with %d peers.", len(peers))
105
if len(peers) == 0 {
106
log.Warning("Running query with no peers!")
@@ -148,7 +148,7 @@ func (r *dhtQueryRunner) Run(peers []*peer.Peer) (*dhtQueryResult, error) {
148
return nil, err
149
}
150
151
-func (r *dhtQueryRunner) addPeerToQuery(next *peer.Peer, benchmark *peer.Peer) {
151
+func (r *dhtQueryRunner) addPeerToQuery(next peer.Peer, benchmark peer.Peer) {
152
if next == nil {
153
// wtf why are peers nil?!?
154
log.Error("Query getting nil peers!!!\n")
@@ -156,7 +156,7 @@ func (r *dhtQueryRunner) addPeerToQuery(next *peer.Peer, benchmark *peer.Peer) {
156
}
157
158
// if new peer further away than whom we got it from, bother (loops)
159
- if benchmark != nil && kb.Closer(benchmark.ID, next.ID, r.query.key) {
159
+ if benchmark != nil && kb.Closer(benchmark.ID(), next.ID(), r.query.key) {
160
return
161
}
162
@@ -200,7 +200,7 @@ func (r *dhtQueryRunner) spawnWorkers() {
200
}
201
}
202
203
-func (r *dhtQueryRunner) queryPeer(p *peer.Peer) {
203
+func (r *dhtQueryRunner) queryPeer(p peer.Peer) {
204
log.Debug("spawned worker for: %v\n", p)
205
206
// make sure we rate limit concurrency.
routing/dht/routing.go
+15
-15
@@ -23,13 +23,13 @@ func (dht *IpfsDHT) PutValue(ctx context.Context, key u.Key, value []byte) error
23
return err
24
}
25
26
- var peers []*peer.Peer
26
+ var peers []peer.Peer
27
for _, route := range dht.routingTables {
28
npeers := route.NearestPeers(kb.ConvertKey(key), KValue)
29
peers = append(peers, npeers...)
30
}
31
32
- query := newQuery(key, func(ctx context.Context, p *peer.Peer) (*dhtQueryResult, error) {
32
+ query := newQuery(key, func(ctx context.Context, p peer.Peer) (*dhtQueryResult, error) {
33
log.Debug("%s PutValue qry part %v", dht.self, p)
34
err := dht.putValueToNetwork(ctx, p, string(key), value)
35
if err != nil {
@@ -65,7 +65,7 @@ func (dht *IpfsDHT) GetValue(ctx context.Context, key u.Key) ([]byte, error) {
65
}
66
67
// setup the Query
68
- query := newQuery(key, func(ctx context.Context, p *peer.Peer) (*dhtQueryResult, error) {
68
+ query := newQuery(key, func(ctx context.Context, p peer.Peer) (*dhtQueryResult, error) {
69
70
val, peers, err := dht.getValueOrPeers(ctx, p, key, routeLevel)
71
if err != nil {
@@ -117,8 +117,8 @@ func (dht *IpfsDHT) Provide(ctx context.Context, key u.Key) error {
117
return nil
118
}
119
120
-func (dht *IpfsDHT) FindProvidersAsync(ctx context.Context, key u.Key, count int) <-chan *peer.Peer {
121
- peerOut := make(chan *peer.Peer, count)
120
+func (dht *IpfsDHT) FindProvidersAsync(ctx context.Context, key u.Key, count int) <-chan peer.Peer {
121
+ peerOut := make(chan peer.Peer, count)
122
go func() {
123
ps := newPeerSet()
124
provs := dht.providers.GetProviders(key)
@@ -136,7 +136,7 @@ func (dht *IpfsDHT) FindProvidersAsync(ctx context.Context, key u.Key, count int
136
peers := dht.routingTables[0].NearestPeers(kb.ConvertKey(key), AlphaValue)
137
for _, pp := range peers {
138
wg.Add(1)
139
- go func(p *peer.Peer) {
139
+ go func(p peer.Peer) {
140
defer wg.Done()
141
pmes, err := dht.findProvidersSingle(ctx, p, key, 0)
142
if err != nil {
@@ -153,7 +153,7 @@ func (dht *IpfsDHT) FindProvidersAsync(ctx context.Context, key u.Key, count int
153
}
154
155
//TODO: this function could also be done asynchronously
156
-func (dht *IpfsDHT) addPeerListAsync(k u.Key, peers []*Message_Peer, ps *peerSet, count int, out chan *peer.Peer) {
156
+func (dht *IpfsDHT) addPeerListAsync(k u.Key, peers []*Message_Peer, ps *peerSet, count int, out chan peer.Peer) {
157
for _, pbp := range peers {
158
159
// construct new peer
@@ -173,7 +173,7 @@ func (dht *IpfsDHT) addPeerListAsync(k u.Key, peers []*Message_Peer, ps *peerSet
173
174
// Find specific Peer
175
// FindPeer searches for a peer with given ID.
176
-func (dht *IpfsDHT) FindPeer(ctx context.Context, id peer.ID) (*peer.Peer, error) {
176
+func (dht *IpfsDHT) FindPeer(ctx context.Context, id peer.ID) (peer.Peer, error) {
177
178
// Check if were already connected to them
179
p, _ := dht.FindLocal(id)
@@ -186,7 +186,7 @@ func (dht *IpfsDHT) FindPeer(ctx context.Context, id peer.ID) (*peer.Peer, error
186
if p == nil {
187
return nil, nil
188
}
189
- if p.ID.Equal(id) {
189
+ if p.ID().Equal(id) {
190
return p, nil
191
}
192
@@ -205,7 +205,7 @@ func (dht *IpfsDHT) FindPeer(ctx context.Context, id peer.ID) (*peer.Peer, error
205
continue
206
}
207
208
- if nxtPeer.ID.Equal(id) {
208
+ if nxtPeer.ID().Equal(id) {
209
return nxtPeer, nil
210
}
211
@@ -214,7 +214,7 @@ func (dht *IpfsDHT) FindPeer(ctx context.Context, id peer.ID) (*peer.Peer, error
214
return nil, u.ErrNotFound
215
}
216
217
-func (dht *IpfsDHT) findPeerMultiple(ctx context.Context, id peer.ID) (*peer.Peer, error) {
217
+func (dht *IpfsDHT) findPeerMultiple(ctx context.Context, id peer.ID) (peer.Peer, error) {
218
219
// Check if were already connected to them
220
p, _ := dht.FindLocal(id)
@@ -230,7 +230,7 @@ func (dht *IpfsDHT) findPeerMultiple(ctx context.Context, id peer.ID) (*peer.Pee
230
}
231
232
// setup query function
233
- query := newQuery(u.Key(id), func(ctx context.Context, p *peer.Peer) (*dhtQueryResult, error) {
233
+ query := newQuery(u.Key(id), func(ctx context.Context, p peer.Peer) (*dhtQueryResult, error) {
234
pmes, err := dht.findPeerSingle(ctx, p, id, routeLevel)
235
if err != nil {
236
log.Error("%s getPeer error: %v", dht.self, err)
@@ -242,7 +242,7 @@ func (dht *IpfsDHT) findPeerMultiple(ctx context.Context, id peer.ID) (*peer.Pee
242
routeLevel++
243
}
244
245
- nxtprs := make([]*peer.Peer, len(plist))
245
+ nxtprs := make([]peer.Peer, len(plist))
246
for i, fp := range plist {
247
nxtp, err := dht.peerFromInfo(fp)
248
if err != nil {
@@ -250,7 +250,7 @@ func (dht *IpfsDHT) findPeerMultiple(ctx context.Context, id peer.ID) (*peer.Pee
250
continue
251
}
252
253
- if nxtp.ID.Equal(id) {
253
+ if nxtp.ID().Equal(id) {
254
return &dhtQueryResult{peer: nxtp, success: true}, nil
255
}
256
@@ -272,7 +272,7 @@ func (dht *IpfsDHT) findPeerMultiple(ctx context.Context, id peer.ID) (*peer.Pee
272
}
273
274
// Ping a peer, log the time it took
275
-func (dht *IpfsDHT) Ping(ctx context.Context, p *peer.Peer) error {
275
+func (dht *IpfsDHT) Ping(ctx context.Context, p peer.Peer) error {
276
// Thoughts: maybe this should accept an ID and do a peer lookup?
277
log.Info("ping %s start", p)
278
routing/dht/util.go
+7
-7
@@ -52,15 +52,15 @@ func newPeerSet() *peerSet {
52
return ps
53
}
54
55
-func (ps *peerSet) Add(p *peer.Peer) {
55
+func (ps *peerSet) Add(p peer.Peer) {
56
ps.lk.Lock()
57
- ps.ps[string(p.ID)] = true
57
+ ps.ps[string(p.ID())] = true
58
ps.lk.Unlock()
59
}
60
61
-func (ps *peerSet) Contains(p *peer.Peer) bool {
61
+func (ps *peerSet) Contains(p peer.Peer) bool {
62
ps.lk.RLock()
63
- _, ok := ps.ps[string(p.ID)]
63
+ _, ok := ps.ps[string(p.ID())]
64
ps.lk.RUnlock()
65
return ok
66
}
@@ -71,12 +71,12 @@ func (ps *peerSet) Size() int {
71
return len(ps.ps)
72
}
73
74
-func (ps *peerSet) AddIfSmallerThan(p *peer.Peer, maxsize int) bool {
74
+func (ps *peerSet) AddIfSmallerThan(p peer.Peer, maxsize int) bool {
75
var success bool
76
ps.lk.Lock()
77
- if _, ok := ps.ps[string(p.ID)]; !ok && len(ps.ps) < maxsize {
77
+ if _, ok := ps.ps[string(p.ID())]; !ok && len(ps.ps) < maxsize {
78
success = true
79
- ps.ps[string(p.ID)] = true
79
+ ps.ps[string(p.ID())] = true
80
}
81
ps.lk.Unlock()
82
return success
routing/kbucket/bucket.go
+5
-5
@@ -23,7 +23,7 @@ func (b *Bucket) find(id peer.ID) *list.Element {
23
b.lk.RLock()
24
defer b.lk.RUnlock()
25
for e := b.list.Front(); e != nil; e = e.Next() {
26
- if e.Value.(*peer.Peer).ID.Equal(id) {
26
+ if e.Value.(peer.Peer).ID().Equal(id) {
27
return e
28
}
29
}
@@ -36,18 +36,18 @@ func (b *Bucket) moveToFront(e *list.Element) {
36
b.lk.Unlock()
37
}
38
39
-func (b *Bucket) pushFront(p *peer.Peer) {
39
+func (b *Bucket) pushFront(p peer.Peer) {
40
b.lk.Lock()
41
b.list.PushFront(p)
42
b.lk.Unlock()
43
}
44
45
-func (b *Bucket) popBack() *peer.Peer {
45
+func (b *Bucket) popBack() peer.Peer {
46
b.lk.Lock()
47
defer b.lk.Unlock()
48
last := b.list.Back()
49
b.list.Remove(last)
50
- return last.Value.(*peer.Peer)
50
+ return last.Value.(peer.Peer)
51
}
52
53
func (b *Bucket) len() int {
@@ -68,7 +68,7 @@ func (b *Bucket) Split(cpl int, target ID) *Bucket {
68
newbuck.list = out
69
e := b.list.Front()
70
for e != nil {
71
- peerID := ConvertPeerID(e.Value.(*peer.Peer).ID)
71
+ peerID := ConvertPeerID(e.Value.(peer.Peer).ID())
72
peerCPL := commonPrefixLen(peerID, target)
73
if peerCPL > cpl {
74
cur := e
routing/kbucket/table.go
+15
-15
@@ -42,10 +42,10 @@ func NewRoutingTable(bucketsize int, localID ID, latency time.Duration) *Routing
42
43
// Update adds or moves the given peer to the front of its respective bucket
44
// If a peer gets removed from a bucket, it is returned
45
-func (rt *RoutingTable) Update(p *peer.Peer) *peer.Peer {
45
+func (rt *RoutingTable) Update(p peer.Peer) peer.Peer {
46
rt.tabLock.Lock()
47
defer rt.tabLock.Unlock()
48
- peerID := ConvertPeerID(p.ID)
48
+ peerID := ConvertPeerID(p.ID())
49
cpl := commonPrefixLen(peerID, rt.local)
50
51
bucketID := cpl
@@ -54,7 +54,7 @@ func (rt *RoutingTable) Update(p *peer.Peer) *peer.Peer {
54
}
55
56
bucket := rt.Buckets[bucketID]
57
- e := bucket.find(p.ID)
57
+ e := bucket.find(p.ID())
58
if e == nil {
59
// New peer, add to bucket
60
if p.GetLatency() > rt.maxLatency {
@@ -93,7 +93,7 @@ func (rt *RoutingTable) Update(p *peer.Peer) *peer.Peer {
93
94
// A helper struct to sort peers by their distance to the local node
95
type peerDistance struct {
96
- p *peer.Peer
96
+ p peer.Peer
97
distance ID
98
}
99
@@ -110,8 +110,8 @@ func (p peerSorterArr) Less(a, b int) bool {
110
111
func copyPeersFromList(target ID, peerArr peerSorterArr, peerList *list.List) peerSorterArr {
112
for e := peerList.Front(); e != nil; e = e.Next() {
113
- p := e.Value.(*peer.Peer)
114
- pID := ConvertPeerID(p.ID)
113
+ p := e.Value.(peer.Peer)
114
+ pID := ConvertPeerID(p.ID())
115
pd := peerDistance{
116
p: p,
117
distance: xor(target, pID),
@@ -126,16 +126,16 @@ func copyPeersFromList(target ID, peerArr peerSorterArr, peerList *list.List) pe
126
}
127
128
// Find a specific peer by ID or return nil
129
-func (rt *RoutingTable) Find(id peer.ID) *peer.Peer {
129
+func (rt *RoutingTable) Find(id peer.ID) peer.Peer {
130
srch := rt.NearestPeers(ConvertPeerID(id), 1)
131
- if len(srch) == 0 || !srch[0].ID.Equal(id) {
131
+ if len(srch) == 0 || !srch[0].ID().Equal(id) {
132
return nil
133
}
134
return srch[0]
135
}
136
137
// NearestPeer returns a single peer that is nearest to the given ID
138
-func (rt *RoutingTable) NearestPeer(id ID) *peer.Peer {
138
+func (rt *RoutingTable) NearestPeer(id ID) peer.Peer {
139
peers := rt.NearestPeers(id, 1)
140
if len(peers) > 0 {
141
return peers[0]
@@ -146,7 +146,7 @@ func (rt *RoutingTable) NearestPeer(id ID) *peer.Peer {
146
}
147
148
// NearestPeers returns a list of the 'count' closest peers to the given ID
149
-func (rt *RoutingTable) NearestPeers(id ID, count int) []*peer.Peer {
149
+func (rt *RoutingTable) NearestPeers(id ID, count int) []peer.Peer {
150
rt.tabLock.RLock()
151
defer rt.tabLock.RUnlock()
152
cpl := commonPrefixLen(id, rt.local)
@@ -178,7 +178,7 @@ func (rt *RoutingTable) NearestPeers(id ID, count int) []*peer.Peer {
178
// Sort by distance to local peer
179
sort.Sort(peerArr)
180
181
- var out []*peer.Peer
181
+ var out []peer.Peer
182
for i := 0; i < count && i < peerArr.Len(); i++ {
183
out = append(out, peerArr[i].p)
184
}
@@ -197,11 +197,11 @@ func (rt *RoutingTable) Size() int {
197
198
// ListPeers takes a RoutingTable and returns a list of all peers from all buckets in the table.
199
// NOTE: This is potentially unsafe... use at your own risk
200
-func (rt *RoutingTable) ListPeers() []*peer.Peer {
201
- var peers []*peer.Peer
200
+func (rt *RoutingTable) ListPeers() []peer.Peer {
201
+ var peers []peer.Peer
202
for _, buck := range rt.Buckets {
203
for e := buck.getIter(); e != nil; e = e.Next() {
204
- peers = append(peers, e.Value.(*peer.Peer))
204
+ peers = append(peers, e.Value.(peer.Peer))
205
}
206
}
207
return peers
@@ -213,6 +213,6 @@ func (rt *RoutingTable) Print() {
213
rt.tabLock.RLock()
214
peers := rt.ListPeers()
215
for i, p := range peers {
216
- fmt.Printf("%d) %s %s\n", i, p.ID.Pretty(), p.GetLatency().String())
216
+ fmt.Printf("%d) %s %s\n", i, p.ID().Pretty(), p.GetLatency().String())
217
}
218
}
routing/kbucket/table_test.go
+24
-25
@@ -10,11 +10,10 @@ import (
10
peer "github.com/jbenet/go-ipfs/peer"
11
)
12
13
-func _randPeer() *peer.Peer {
14
- p := new(peer.Peer)
15
- p.ID = make(peer.ID, 16)
16
- crand.Read(p.ID)
17
- return p
13
+func _randPeer() peer.Peer {
14
+ id := make(peer.ID, 16)
15
+ crand.Read(id)
16
+ return peer.WithID(id)
17
}
18
19
func _randID() ID {
@@ -29,25 +28,25 @@ func _randID() ID {
28
func TestBucket(t *testing.T) {
29
b := newBucket()
30
32
- peers := make([]*peer.Peer, 100)
31
+ peers := make([]peer.Peer, 100)
32
for i := 0; i < 100; i++ {
33
peers[i] = _randPeer()
34
b.pushFront(peers[i])
35
}
36
37
local := _randPeer()
39
- localID := ConvertPeerID(local.ID)
38
+ localID := ConvertPeerID(local.ID())
39
40
i := rand.Intn(len(peers))
42
- e := b.find(peers[i].ID)
41
+ e := b.find(peers[i].ID())
42
if e == nil {
43
t.Errorf("Failed to find peer: %v", peers[i])
44
}
45
47
- spl := b.Split(0, ConvertPeerID(local.ID))
46
+ spl := b.Split(0, ConvertPeerID(local.ID()))
47
llist := b.list
48
for e := llist.Front(); e != nil; e = e.Next() {
50
- p := ConvertPeerID(e.Value.(*peer.Peer).ID)
49
+ p := ConvertPeerID(e.Value.(peer.Peer).ID())
50
cpl := commonPrefixLen(p, localID)
51
if cpl > 0 {
52
t.Fatalf("Split failed. found id with cpl > 0 in 0 bucket")
@@ -56,7 +55,7 @@ func TestBucket(t *testing.T) {
55
56
rlist := spl.list
57
for e := rlist.Front(); e != nil; e = e.Next() {
59
- p := ConvertPeerID(e.Value.(*peer.Peer).ID)
58
+ p := ConvertPeerID(e.Value.(peer.Peer).ID())
59
cpl := commonPrefixLen(p, localID)
60
if cpl == 0 {
61
t.Fatalf("Split failed. found id with cpl == 0 in non 0 bucket")
@@ -67,9 +66,9 @@ func TestBucket(t *testing.T) {
66
// Right now, this just makes sure that it doesnt hang or crash
67
func TestTableUpdate(t *testing.T) {
68
local := _randPeer()
70
- rt := NewRoutingTable(10, ConvertPeerID(local.ID), time.Hour)
69
+ rt := NewRoutingTable(10, ConvertPeerID(local.ID()), time.Hour)
70
72
- peers := make([]*peer.Peer, 100)
71
+ peers := make([]peer.Peer, 100)
72
for i := 0; i < 100; i++ {
73
peers[i] = _randPeer()
74
}
@@ -93,33 +92,33 @@ func TestTableUpdate(t *testing.T) {
92
93
func TestTableFind(t *testing.T) {
94
local := _randPeer()
96
- rt := NewRoutingTable(10, ConvertPeerID(local.ID), time.Hour)
95
+ rt := NewRoutingTable(10, ConvertPeerID(local.ID()), time.Hour)
96
98
- peers := make([]*peer.Peer, 100)
97
+ peers := make([]peer.Peer, 100)
98
for i := 0; i < 5; i++ {
99
peers[i] = _randPeer()
100
rt.Update(peers[i])
101
}
102
103
t.Logf("Searching for peer: '%s'", peers[2])
105
- found := rt.NearestPeer(ConvertPeerID(peers[2].ID))
106
- if !found.ID.Equal(peers[2].ID) {
104
+ found := rt.NearestPeer(ConvertPeerID(peers[2].ID()))
105
+ if !found.ID().Equal(peers[2].ID()) {
106
t.Fatalf("Failed to lookup known node...")
107
}
108
}
109
110
func TestTableFindMultiple(t *testing.T) {
111
local := _randPeer()
113
- rt := NewRoutingTable(20, ConvertPeerID(local.ID), time.Hour)
112
+ rt := NewRoutingTable(20, ConvertPeerID(local.ID()), time.Hour)
113
115
- peers := make([]*peer.Peer, 100)
114
+ peers := make([]peer.Peer, 100)
115
for i := 0; i < 18; i++ {
116
peers[i] = _randPeer()
117
rt.Update(peers[i])
118
}
119
120
t.Logf("Searching for peer: '%s'", peers[2])
122
- found := rt.NearestPeers(ConvertPeerID(peers[2].ID), 15)
121
+ found := rt.NearestPeers(ConvertPeerID(peers[2].ID()), 15)
122
if len(found) != 15 {
123
t.Fatalf("Got back different number of peers than we expected.")
124
}
@@ -131,7 +130,7 @@ func TestTableFindMultiple(t *testing.T) {
130
func TestTableMultithreaded(t *testing.T) {
131
local := peer.ID("localPeer")
132
tab := NewRoutingTable(20, ConvertPeerID(local), time.Hour)
134
- var peers []*peer.Peer
133
+ var peers []peer.Peer
134
for i := 0; i < 500; i++ {
135
peers = append(peers, _randPeer())
136
}
@@ -156,7 +155,7 @@ func TestTableMultithreaded(t *testing.T) {
155
go func() {
156
for i := 0; i < 1000; i++ {
157
n := rand.Intn(len(peers))
159
- tab.Find(peers[n].ID)
158
+ tab.Find(peers[n].ID())
159
}
160
done <- struct{}{}
161
}()
@@ -170,7 +169,7 @@ func BenchmarkUpdates(b *testing.B) {
169
local := ConvertKey("localKey")
170
tab := NewRoutingTable(20, local, time.Hour)
171
173
- var peers []*peer.Peer
172
+ var peers []peer.Peer
173
for i := 0; i < b.N; i++ {
174
peers = append(peers, _randPeer())
175
}
@@ -186,7 +185,7 @@ func BenchmarkFinds(b *testing.B) {
185
local := ConvertKey("localKey")
186
tab := NewRoutingTable(20, local, time.Hour)
187
189
- var peers []*peer.Peer
188
+ var peers []peer.Peer
189
for i := 0; i < b.N; i++ {
190
peers = append(peers, _randPeer())
191
tab.Update(peers[i])
@@ -194,6 +193,6 @@ func BenchmarkFinds(b *testing.B) {
193
194
b.StartTimer()
195
for i := 0; i < b.N; i++ {
197
- tab.Find(peers[i].ID)
196
+ tab.Find(peers[i].ID())
197
}
198
}
routing/mock/routing.go
+13
-13
@@ -17,10 +17,10 @@ var _ routing.IpfsRouting = &MockRouter{}
17
type MockRouter struct {
18
datastore ds.Datastore
19
hashTable RoutingServer
20
- peer *peer.Peer
20
+ peer peer.Peer
21
}
22
23
-func NewMockRouter(local *peer.Peer, dstore ds.Datastore) routing.IpfsRouting {
23
+func NewMockRouter(local peer.Peer, dstore ds.Datastore) routing.IpfsRouting {
24
return &MockRouter{
25
datastore: dstore,
26
peer: local,
@@ -50,16 +50,16 @@ func (mr *MockRouter) GetValue(ctx context.Context, key u.Key) ([]byte, error) {
50
return data, nil
51
}
52
53
-func (mr *MockRouter) FindProviders(ctx context.Context, key u.Key) ([]*peer.Peer, error) {
53
+func (mr *MockRouter) FindProviders(ctx context.Context, key u.Key) ([]peer.Peer, error) {
54
return nil, nil
55
}
56
57
-func (mr *MockRouter) FindPeer(ctx context.Context, pid peer.ID) (*peer.Peer, error) {
57
+func (mr *MockRouter) FindPeer(ctx context.Context, pid peer.ID) (peer.Peer, error) {
58
return nil, nil
59
}
60
61
-func (mr *MockRouter) FindProvidersAsync(ctx context.Context, k u.Key, max int) <-chan *peer.Peer {
62
- out := make(chan *peer.Peer)
61
+func (mr *MockRouter) FindProvidersAsync(ctx context.Context, k u.Key, max int) <-chan peer.Peer {
62
+ out := make(chan peer.Peer)
63
go func() {
64
defer close(out)
65
for i, p := range mr.hashTable.Providers(k) {
@@ -81,11 +81,11 @@ func (mr *MockRouter) Provide(_ context.Context, key u.Key) error {
81
}
82
83
type RoutingServer interface {
84
- Announce(*peer.Peer, u.Key) error
84
+ Announce(peer.Peer, u.Key) error
85
86
- Providers(u.Key) []*peer.Peer
86
+ Providers(u.Key) []peer.Peer
87
88
- Client(p *peer.Peer) routing.IpfsRouting
88
+ Client(p peer.Peer) routing.IpfsRouting
89
}
90
91
func VirtualRoutingServer() RoutingServer {
@@ -99,7 +99,7 @@ type hashTable struct {
99
providers map[u.Key]peer.Map
100
}
101
102
-func (rs *hashTable) Announce(p *peer.Peer, k u.Key) error {
102
+func (rs *hashTable) Announce(p peer.Peer, k u.Key) error {
103
rs.lock.Lock()
104
defer rs.lock.Unlock()
105
@@ -111,10 +111,10 @@ func (rs *hashTable) Announce(p *peer.Peer, k u.Key) error {
111
return nil
112
}
113
114
-func (rs *hashTable) Providers(k u.Key) []*peer.Peer {
114
+func (rs *hashTable) Providers(k u.Key) []peer.Peer {
115
rs.lock.RLock()
116
defer rs.lock.RUnlock()
117
- ret := make([]*peer.Peer, 0)
117
+ ret := make([]peer.Peer, 0)
118
peerset, ok := rs.providers[k]
119
if !ok {
120
return ret
@@ -131,7 +131,7 @@ func (rs *hashTable) Providers(k u.Key) []*peer.Peer {
131
return ret
132
}
133
134
-func (rs *hashTable) Client(p *peer.Peer) routing.IpfsRouting {
134
+func (rs *hashTable) Client(p peer.Peer) routing.IpfsRouting {
135
return &MockRouter{
136
peer: p,
137
hashTable: rs,
routing/mock/routing_test.go
+9
-15
@@ -20,9 +20,7 @@ func TestKeyNotFound(t *testing.T) {
20
21
func TestSetAndGet(t *testing.T) {
22
pid := peer.ID([]byte("the peer id"))
23
- p := &peer.Peer{
24
- ID: pid,
25
- }
23
+ p := peer.WithID(pid)
24
k := u.Key("42")
25
rs := VirtualRoutingServer()
26
err := rs.Announce(p, k)
@@ -34,7 +32,7 @@ func TestSetAndGet(t *testing.T) {
32
t.Fatal("should be one")
33
}
34
for _, elem := range providers {
37
- if bytes.Equal(elem.ID, pid) {
35
+ if bytes.Equal(elem.ID(), pid) {
36
return
37
}
38
}
@@ -42,7 +40,7 @@ func TestSetAndGet(t *testing.T) {
40
}
41
42
func TestClientFindProviders(t *testing.T) {
45
- peer := &peer.Peer{ID: []byte("42")}
43
+ peer := peer.WithIDString("42")
44
rs := VirtualRoutingServer()
45
client := rs.Client(peer)
46
@@ -57,7 +55,7 @@ func TestClientFindProviders(t *testing.T) {
55
56
isInHT := false
57
for _, p := range providersFromHashTable {
60
- if bytes.Equal(p.ID, peer.ID) {
58
+ if bytes.Equal(p.ID(), peer.ID()) {
59
isInHT = true
60
}
61
}
@@ -67,7 +65,7 @@ func TestClientFindProviders(t *testing.T) {
65
providersFromClient := client.FindProvidersAsync(context.Background(), u.Key("hello"), max)
66
isInClient := false
67
for p := range providersFromClient {
70
- if bytes.Equal(p.ID, peer.ID) {
68
+ if bytes.Equal(p.ID(), peer.ID()) {
69
isInClient = true
70
}
71
}
@@ -81,9 +79,7 @@ func TestClientOverMax(t *testing.T) {
79
k := u.Key("hello")
80
numProvidersForHelloKey := 100
81
for i := 0; i < numProvidersForHelloKey; i++ {
84
- peer := &peer.Peer{
85
- ID: []byte(string(i)),
86
- }
82
+ peer := peer.WithIDString(string(i))
83
err := rs.Announce(peer, k)
84
if err != nil {
85
t.Fatal(err)
@@ -96,7 +92,7 @@ func TestClientOverMax(t *testing.T) {
92
}
93
94
max := 10
99
- peer := &peer.Peer{ID: []byte("TODO")}
95
+ peer := peer.WithIDString("TODO")
96
client := rs.Client(peer)
97
98
providersFromClient := client.FindProvidersAsync(context.Background(), k, max)
@@ -118,9 +114,7 @@ func TestCanceledContext(t *testing.T) {
114
i := 0
115
go func() { // infinite stream
116
for {
121
- peer := &peer.Peer{
122
- ID: []byte(string(i)),
123
- }
117
+ peer := peer.WithIDString(string(i))
118
err := rs.Announce(peer, k)
119
if err != nil {
120
t.Fatal(err)
@@ -129,7 +123,7 @@ func TestCanceledContext(t *testing.T) {
123
}
124
}()
125
132
- local := &peer.Peer{ID: []byte("peer id doesn't matter")}
126
+ local := peer.WithIDString("peer id doesn't matter")
127
client := rs.Client(local)
128
129
t.Log("warning: max is finite so this test is non-deterministic")
routing/routing.go
+2
-2
@@ -10,7 +10,7 @@ import (
10
// IpfsRouting is the routing module interface
11
// It is implemented by things like DHTs, etc.
12
type IpfsRouting interface {
13
- FindProvidersAsync(context.Context, u.Key, int) <-chan *peer.Peer
13
+ FindProvidersAsync(context.Context, u.Key, int) <-chan peer.Peer
14
15
// Basic Put/Get
16
@@ -28,5 +28,5 @@ type IpfsRouting interface {
28
29
// Find specific Peer
30
// FindPeer searches for a peer with given ID.
31
- FindPeer(context.Context, peer.ID) (*peer.Peer, error)
31
+ FindPeer(context.Context, peer.ID) (peer.Peer, error)
32
}