-
Notifications
You must be signed in to change notification settings - Fork 282
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Fix #787: Connectivity fixes #792
Changes from all commits
196aa23
99be078
2a2f8c0
d5cef9d
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
Original file line number | Diff line number | Diff line change |
---|---|---|
|
@@ -8,11 +8,10 @@ import ( | |
"net/http" | ||
"time" | ||
|
||
"github.com/ipfs/ipfs-cluster/api" | ||
|
||
p2phttp "github.com/hsanjuan/go-libp2p-http" | ||
libp2p "github.com/libp2p/go-libp2p" | ||
ipnet "github.com/libp2p/go-libp2p-interface-pnet" | ||
peer "github.com/libp2p/go-libp2p-peer" | ||
peerstore "github.com/libp2p/go-libp2p-peerstore" | ||
pnet "github.com/libp2p/go-libp2p-pnet" | ||
madns "github.com/multiformats/go-multiaddr-dns" | ||
|
@@ -41,11 +40,15 @@ func (c *defaultClient) defaultTransport() { | |
func (c *defaultClient) enableLibp2p() error { | ||
c.defaultTransport() | ||
|
||
pid, addr, err := api.Libp2pMultiaddrSplit(c.config.APIAddr) | ||
pinfo, err := peerstore.InfoFromP2pAddr(c.config.APIAddr) | ||
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. A function in |
||
if err != nil { | ||
return err | ||
} | ||
|
||
if len(pinfo.Addrs) == 0 { | ||
return errors.New("APIAddr only includes a Peer ID") | ||
} | ||
|
||
var prot ipnet.Protector | ||
if c.config.ProtectorKey != nil && len(c.config.ProtectorKey) > 0 { | ||
if len(c.config.ProtectorKey) != 32 { | ||
|
@@ -67,16 +70,16 @@ func (c *defaultClient) enableLibp2p() error { | |
|
||
ctx, cancel := context.WithTimeout(c.ctx, ResolveTimeout) | ||
defer cancel() | ||
resolvedAddrs, err := madns.Resolve(ctx, addr) | ||
resolvedAddrs, err := madns.Resolve(ctx, pinfo.Addrs[0]) | ||
if err != nil { | ||
return err | ||
} | ||
|
||
h.Peerstore().AddAddrs(pid, resolvedAddrs, peerstore.PermanentAddrTTL) | ||
h.Peerstore().AddAddrs(pinfo.ID, resolvedAddrs, peerstore.PermanentAddrTTL) | ||
c.transport.RegisterProtocol("libp2p", p2phttp.NewTransport(h)) | ||
c.net = "libp2p" | ||
c.p2p = h | ||
c.hostname = pid.Pretty() | ||
c.hostname = peer.IDB58Encode(pinfo.ID) | ||
return nil | ||
} | ||
|
||
|
Original file line number | Diff line number | Diff line change |
---|---|---|
|
@@ -24,6 +24,7 @@ import ( | |
host "github.com/libp2p/go-libp2p-host" | ||
dht "github.com/libp2p/go-libp2p-kad-dht" | ||
peer "github.com/libp2p/go-libp2p-peer" | ||
peerstore "github.com/libp2p/go-libp2p-peerstore" | ||
ma "github.com/multiformats/go-multiaddr" | ||
|
||
ocgorpc "github.com/lanzafame/go-libp2p-ocgorpc" | ||
|
@@ -36,7 +37,10 @@ import ( | |
// consensus layer. | ||
var ReadyTimeout = 30 * time.Second | ||
|
||
var pingMetricName = "ping" | ||
const ( | ||
pingMetricName = "ping" | ||
bootstrapCount = 3 | ||
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. This is the number of peers we attempt to connect. We block until 3 connections succeed or we run out of addresses to try. After that the DHT will take care of discovering the swarm in the background. |
||
) | ||
|
||
// Cluster is the main IPFS cluster component. It provides | ||
// the go-API for it and orchestrates the components that make up the system. | ||
|
@@ -116,9 +120,7 @@ func NewCluster( | |
|
||
logger.Infof("IPFS Cluster v%s listening on:\n%s\n", version.Version, listenAddrs) | ||
|
||
// Note, we already loaded peers from peerstore into the host | ||
// in daemon.go. | ||
peerManager := pstoremgr.New(host, cfg.GetPeerstorePath()) | ||
peerManager := pstoremgr.New(ctx, host, cfg.GetPeerstorePath()) | ||
|
||
c := &Cluster{ | ||
ctx: ctx, | ||
|
@@ -144,6 +146,18 @@ func NewCluster( | |
readyB: false, | ||
} | ||
|
||
// Import known cluster peers from peerstore file. Set | ||
// a non permanent TTL. | ||
c.peerManager.ImportPeersFromPeerstore(false, peerstore.AddressTTL) | ||
// Attempt to connect to some peers (up to bootstrapCount) | ||
actualCount := c.peerManager.Bootstrap(bootstrapCount) | ||
// We cannot warn about this as this is normal if going to Join() later | ||
logger.Debugf("bootstrap count %d", actualCount) | ||
// Bootstrap the DHT now that we possibly have some connections | ||
c.dht.Bootstrap(c.ctx) | ||
|
||
// After setupRPC components can do their tasks with a fully operative | ||
// routed libp2p host with some connections and a working DHT (hopefully). | ||
err = c.setupRPC() | ||
if err != nil { | ||
c.Shutdown(ctx) | ||
|
@@ -465,9 +479,6 @@ This might be due to one or several causes: | |
|
||
// Cluster is ready. | ||
|
||
// Bootstrap the DHT now that we possibly have some connections | ||
c.dht.Bootstrap(c.ctx) | ||
|
||
peers, err := c.consensus.Peers(ctx) | ||
if err != nil { | ||
logger.Error(err) | ||
|
@@ -632,12 +643,24 @@ func (c *Cluster) ID(ctx context.Context) *api.ID { | |
peers, _ = c.consensus.Peers(ctx) | ||
} | ||
|
||
clusterPeerInfos := c.peerManager.PeerInfos(peers) | ||
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. This change block is because the PeerManager does not return |
||
addresses := []api.Multiaddr{} | ||
for _, pinfo := range clusterPeerInfos { | ||
addrs, err := peerstore.InfoToP2pAddrs(&pinfo) | ||
if err != nil { | ||
continue | ||
} | ||
for _, a := range addrs { | ||
addresses = append(addresses, api.NewMultiaddrWithValue(a)) | ||
} | ||
} | ||
|
||
return &api.ID{ | ||
ID: c.id, | ||
//PublicKey: c.host.Peerstore().PubKey(c.id), | ||
Addresses: addrs, | ||
ClusterPeers: peers, | ||
ClusterPeersAddresses: c.peerManager.PeersAddresses(peers), | ||
ClusterPeersAddresses: addresses, | ||
Version: version.Version.String(), | ||
RPCProtocolVersion: version.RPCProtocol, | ||
IPFS: ipfsID, | ||
|
@@ -720,20 +743,15 @@ func (c *Cluster) Join(ctx context.Context, addr ma.Multiaddr) error { | |
|
||
logger.Debugf("Join(%s)", addr) | ||
|
||
pid, _, err := api.Libp2pMultiaddrSplit(addr) | ||
// Add peer to peerstore so we can talk to it (and connect) | ||
pid, err := c.peerManager.ImportPeer(addr, true, peerstore.PermanentAddrTTL) | ||
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. returning the pid from that function allowed some simplication here, but it does not change behaviour. |
||
if err != nil { | ||
logger.Error(err) | ||
return err | ||
} | ||
|
||
// Bootstrap to myself | ||
if pid == c.id { | ||
return nil | ||
} | ||
|
||
// Add peer to peerstore so we can talk to it (and connect) | ||
c.peerManager.ImportPeer(addr, true) | ||
|
||
// Note that PeerAdd() on the remote peer will | ||
// figure out what our real address is (obviously not | ||
// ListenAddr). | ||
|
Original file line number | Diff line number | Diff line change |
---|---|---|
|
@@ -9,6 +9,7 @@ import ( | |
ipfslite "github.com/hsanjuan/ipfs-lite" | ||
dshelp "github.com/ipfs/go-ipfs-ds-help" | ||
"github.com/ipfs/ipfs-cluster/api" | ||
"github.com/ipfs/ipfs-cluster/pstoremgr" | ||
"github.com/ipfs/ipfs-cluster/state" | ||
"github.com/ipfs/ipfs-cluster/state/dsstate" | ||
multihash "github.com/multiformats/go-multihash" | ||
|
@@ -22,13 +23,15 @@ import ( | |
host "github.com/libp2p/go-libp2p-host" | ||
dht "github.com/libp2p/go-libp2p-kad-dht" | ||
peer "github.com/libp2p/go-libp2p-peer" | ||
peerstore "github.com/libp2p/go-libp2p-peerstore" | ||
pubsub "github.com/libp2p/go-libp2p-pubsub" | ||
) | ||
|
||
var logger = logging.Logger("crdt") | ||
|
||
var ( | ||
blocksNs = "b" // blockstore namespace | ||
blocksNs = "b" // blockstore namespace | ||
connMgrTag = "crdt" | ||
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. The connmanager needs a tag when we tag something for Protect(). |
||
) | ||
|
||
// Common variables for the module. | ||
|
@@ -48,7 +51,8 @@ type Consensus struct { | |
|
||
trustedPeers sync.Map | ||
|
||
host host.Host | ||
host host.Host | ||
peerManager *pstoremgr.Manager | ||
|
||
store ds.Datastore | ||
namespace ds.Key | ||
|
@@ -85,21 +89,17 @@ func New( | |
ctx, cancel := context.WithCancel(context.Background()) | ||
|
||
css := &Consensus{ | ||
ctx: ctx, | ||
cancel: cancel, | ||
config: cfg, | ||
host: host, | ||
dht: dht, | ||
store: store, | ||
namespace: ds.NewKey(cfg.DatastoreNamespace), | ||
pubsub: pubsub, | ||
rpcReady: make(chan struct{}, 1), | ||
readyCh: make(chan struct{}, 1), | ||
} | ||
|
||
// Set up a fast-lookup trusted peers cache. | ||
for _, p := range css.config.TrustedPeers { | ||
css.Trust(ctx, p) | ||
ctx: ctx, | ||
cancel: cancel, | ||
config: cfg, | ||
host: host, | ||
peerManager: pstoremgr.New(ctx, host, ""), | ||
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. We keep a peerManager around to set higher priority to trusted peers |
||
dht: dht, | ||
store: store, | ||
namespace: ds.NewKey(cfg.DatastoreNamespace), | ||
pubsub: pubsub, | ||
rpcReady: make(chan struct{}, 1), | ||
readyCh: make(chan struct{}, 1), | ||
} | ||
|
||
go css.setup() | ||
|
@@ -113,6 +113,12 @@ func (css *Consensus) setup() { | |
case <-css.rpcReady: | ||
} | ||
|
||
// Set up a fast-lookup trusted peers cache. | ||
// Protect these peers in the ConnMgr | ||
for _, p := range css.config.TrustedPeers { | ||
css.Trust(css.ctx, p) | ||
} | ||
|
||
// Hash the cluster name and produce the topic name from there | ||
// as a way to avoid pubsub topic collisions with other | ||
// pubsub applications potentially when both potentially use | ||
|
@@ -296,9 +302,18 @@ func (css *Consensus) IsTrustedPeer(ctx context.Context, pid peer.ID) bool { | |
return ok | ||
} | ||
|
||
// Trust marks a peer as "trusted". | ||
// Trust marks a peer as "trusted". It makes sure it is trusted as issuer | ||
// for pubsub updates, it is protected in the connection manager, it | ||
// has the highest priority when the peerstore is saved, and it's addresses | ||
// are always remembered. | ||
func (css *Consensus) Trust(ctx context.Context, pid peer.ID) error { | ||
css.trustedPeers.Store(pid, struct{}{}) | ||
if conman := css.host.ConnManager(); conman != nil { | ||
conman.Protect(pid, connMgrTag) | ||
} | ||
css.peerManager.SetPriority(pid, 0) | ||
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. 0 means highest priority. |
||
addrs := css.host.Peerstore().Addrs(pid) | ||
css.host.Peerstore().SetAddrs(pid, addrs, peerstore.PermanentAddrTTL) | ||
return nil | ||
} | ||
|
||
|
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
These two changes are two fix a broken tests that I must have missed when updating deps.