2017-01-30 12:12:25 +00:00
|
|
|
package ipfscluster
|
|
|
|
|
|
|
|
import (
|
2017-11-24 14:12:47 +00:00
|
|
|
"context"
|
2017-11-08 19:04:04 +00:00
|
|
|
"fmt"
|
2017-11-24 14:12:47 +00:00
|
|
|
"time"
|
2017-11-14 11:26:42 +00:00
|
|
|
|
2017-11-08 19:04:04 +00:00
|
|
|
host "github.com/libp2p/go-libp2p-host"
|
2017-01-30 12:12:25 +00:00
|
|
|
peer "github.com/libp2p/go-libp2p-peer"
|
|
|
|
peerstore "github.com/libp2p/go-libp2p-peerstore"
|
|
|
|
ma "github.com/multiformats/go-multiaddr"
|
2017-11-24 14:12:47 +00:00
|
|
|
madns "github.com/multiformats/go-multiaddr-dns"
|
2017-01-30 12:12:25 +00:00
|
|
|
)
|
|
|
|
|
2017-11-08 19:04:04 +00:00
|
|
|
// peerManager provides wrappers peerset control
|
2017-01-30 12:12:25 +00:00
|
|
|
type peerManager struct {
|
2017-11-08 19:04:04 +00:00
|
|
|
host host.Host
|
2018-01-16 19:57:54 +00:00
|
|
|
ctx context.Context
|
2017-01-30 12:12:25 +00:00
|
|
|
}
|
|
|
|
|
2017-11-08 19:04:04 +00:00
|
|
|
func newPeerManager(h host.Host) *peerManager {
|
2018-01-16 19:57:54 +00:00
|
|
|
return &peerManager{
|
|
|
|
ctx: context.Background(),
|
|
|
|
host: h,
|
|
|
|
}
|
2017-01-30 12:12:25 +00:00
|
|
|
}
|
|
|
|
|
2018-01-16 19:57:54 +00:00
|
|
|
func (pm *peerManager) addPeer(addr ma.Multiaddr, connect bool) error {
|
2017-11-08 19:04:04 +00:00
|
|
|
logger.Debugf("adding peer address %s", addr)
|
2017-02-02 22:52:06 +00:00
|
|
|
pid, decapAddr, err := multiaddrSplit(addr)
|
2017-01-30 12:12:25 +00:00
|
|
|
if err != nil {
|
2017-02-02 22:52:06 +00:00
|
|
|
return err
|
2017-01-30 12:12:25 +00:00
|
|
|
}
|
2017-11-08 19:04:04 +00:00
|
|
|
pm.host.Peerstore().AddAddr(pid, decapAddr, peerstore.PermanentAddrTTL)
|
2017-11-24 14:12:47 +00:00
|
|
|
|
|
|
|
// dns multiaddresses need to be resolved because libp2p only does that
|
|
|
|
// on explicit bhost.Connect().
|
|
|
|
if madns.Matches(addr) {
|
|
|
|
ctx, cancel := context.WithTimeout(context.Background(), time.Second*2)
|
|
|
|
defer cancel()
|
|
|
|
resolvedAddrs, err := madns.Resolve(ctx, addr)
|
|
|
|
if err != nil {
|
|
|
|
logger.Error(err)
|
|
|
|
return err
|
|
|
|
}
|
2018-01-16 19:57:54 +00:00
|
|
|
pm.importAddresses(resolvedAddrs, connect)
|
|
|
|
}
|
|
|
|
if connect {
|
|
|
|
pm.host.Network().DialPeer(pm.ctx, pid)
|
2017-11-24 14:12:47 +00:00
|
|
|
}
|
2017-02-02 22:52:06 +00:00
|
|
|
return nil
|
2017-01-30 12:12:25 +00:00
|
|
|
}
|
|
|
|
|
2017-11-08 19:04:04 +00:00
|
|
|
func (pm *peerManager) rmPeer(pid peer.ID) error {
|
|
|
|
logger.Debugf("forgetting peer %s", pid.Pretty())
|
|
|
|
pm.host.Peerstore().ClearAddrs(pid)
|
2017-01-30 12:12:25 +00:00
|
|
|
return nil
|
|
|
|
}
|
|
|
|
|
2017-02-02 22:52:06 +00:00
|
|
|
// cluster peer addresses (NOT including ourselves)
|
2017-11-08 19:04:04 +00:00
|
|
|
func (pm *peerManager) addresses(peers []peer.ID) []ma.Multiaddr {
|
2017-10-27 20:11:14 +00:00
|
|
|
addrs := []ma.Multiaddr{}
|
2017-11-08 19:04:04 +00:00
|
|
|
if peers == nil {
|
|
|
|
return addrs
|
|
|
|
}
|
|
|
|
|
|
|
|
for _, p := range peers {
|
|
|
|
if p == pm.host.ID() {
|
|
|
|
continue
|
|
|
|
}
|
|
|
|
peerAddr, _ := ma.NewMultiaddr(fmt.Sprintf("/ipfs/%s", peer.IDB58Encode(p)))
|
|
|
|
for _, a := range pm.host.Peerstore().Addrs(p) {
|
|
|
|
addrs = append(addrs, a.Encapsulate(peerAddr))
|
2017-02-02 22:52:06 +00:00
|
|
|
}
|
2017-01-30 12:12:25 +00:00
|
|
|
}
|
2017-02-02 22:52:06 +00:00
|
|
|
return addrs
|
|
|
|
}
|
2017-01-30 12:12:25 +00:00
|
|
|
|
2018-01-16 19:57:54 +00:00
|
|
|
func (pm *peerManager) importAddresses(addrs []ma.Multiaddr, connect bool) error {
|
2017-11-08 19:04:04 +00:00
|
|
|
for _, a := range addrs {
|
2018-01-16 19:57:54 +00:00
|
|
|
pm.addPeer(a, connect)
|
2017-10-31 10:20:14 +00:00
|
|
|
}
|
2017-01-30 12:12:25 +00:00
|
|
|
return nil
|
|
|
|
}
|