@cryptotaxi247 / kubo / commits / 175da4f58

feat(core) supervise bootstrap connections

License: MIT Signed-off-by: Brian Tiger Chow <brian@perfmode.com>

Brian Tiger Chow committed Dec 8, 2014 at 05:11 UTC 175da4f584a1790afdaa9d0d237643ac6c3c675d
2 files changed +112 -37
core/bootstrap.go new
+110
@@ -0,0 +1,110 @@
1 +package core
2 +
3 +import (
4 + "sync"
5 + "time"
6 +
7 + context "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
8 + ma "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-multiaddr"
9 + "github.com/jbenet/go-ipfs/config"
10 + inet "github.com/jbenet/go-ipfs/net"
11 + "github.com/jbenet/go-ipfs/peer"
12 + "github.com/jbenet/go-ipfs/routing/dht"
13 +)
14 +
15 +const period time.Duration = 30 * time.Second
16 +const timeout time.Duration = period / 3
17 +
18 +func superviseConnections(parent context.Context,
19 + n *IpfsNode,
20 + route *dht.IpfsDHT,
21 + store peer.Peerstore,
22 + peers []*config.BootstrapPeer) error {
23 +
24 + for {
25 + ctx, _ := context.WithTimeout(parent, timeout)
26 + // TODO get config from disk so |peers| always reflects the latest
27 + // information
28 + if err := bootstrap(ctx, n.Network, route, store, peers); err != nil {
29 + log.Error(err)
30 + }
31 + select {
32 + case <-parent.Done():
33 + return parent.Err()
34 + case <-time.Tick(period):
35 + }
36 + }
37 + return nil
38 +}
39 +
40 +func bootstrap(ctx context.Context,
41 + n inet.Network,
42 + r *dht.IpfsDHT,
43 + ps peer.Peerstore,
44 + boots []*config.BootstrapPeer) error {
45 +
46 + var peers []peer.Peer
47 + for _, bootstrap := range boots {
48 + p, err := toPeer(ps, bootstrap)
49 + if err != nil {
50 + return err
51 + }
52 + peers = append(peers, p)
53 + }
54 +
55 + var notConnected []peer.Peer
56 + for _, p := range peers {
57 + if !n.IsConnected(p) {
58 + notConnected = append(notConnected, p)
59 + }
60 + }
61 + for _, p := range notConnected {
62 + log.Infof("not connected to %v", p)
63 + }
64 + if err := connect(ctx, r, notConnected); err != nil {
65 + return err
66 + }
67 + return nil
68 +}
69 +
70 +func connect(ctx context.Context, r *dht.IpfsDHT, peers []peer.Peer) error {
71 + var wg sync.WaitGroup
72 + for _, p := range peers {
73 +
74 + // performed asynchronously because when performed synchronously, if
75 + // one `Connect` call hangs, subsequent calls are more likely to
76 + // fail/abort due to an expiring context.
77 +
78 + wg.Add(1)
79 + go func(p peer.Peer) {
80 + defer wg.Done()
81 + err := r.Connect(ctx, p)
82 + if err != nil {
83 + log.Event(ctx, "bootstrapFailed", p)
84 + log.Criticalf("failed to bootstrap with %v", p)
85 + return
86 + }
87 + log.Event(ctx, "bootstrapSuccess", p)
88 + log.Infof("bootstrapped with %v", p)
89 + }(p)
90 + }
91 + wg.Wait()
92 + return nil
93 +}
94 +
95 +func toPeer(ps peer.Peerstore, bootstrap *config.BootstrapPeer) (peer.Peer, error) {
96 + id, err := peer.DecodePrettyID(bootstrap.PeerID)
97 + if err != nil {
98 + return nil, err
99 + }
100 + p, err := ps.FindOrCreate(id)
101 + if err != nil {
102 + return nil, err
103 + }
104 + maddr, err := ma.NewMultiaddr(bootstrap.Address)
105 + if err != nil {
106 + return nil, err
107 + }
108 + p.AddAddress(maddr)
109 + return p, nil
110 +}
core/core.go
+2 -37
@@ -179,13 +179,13 @@ func NewIpfsNode(cfg *config.Config, online bool) (n *IpfsNode, err error) {
179
180 n.Exchange = bitswap.New(ctx, n.Identity, bitswapNetwork, n.Routing, blockstore, alwaysSendToPeer)
181
182 - // TODO consider connection supervision into the Network. We've
182 + // TODO consider moving connection supervision into the Network. We've
183 // discussed improvements to this Node constructor. One improvement
184 // would be to make the node configurable, allowing clients to inject
185 // an Exchange, Network, or Routing component and have the constructor
186 // manage the wiring. In that scenario, this dangling function is a bit
187 // awkward.
188 - go initConnections(ctx, n.Config.Bootstrap, n.Peerstore, dhtRouting)
188 + go superviseConnections(ctx, n, dhtRouting, n.Peerstore, n.Config.Bootstrap)
189 }
190
191 // TODO(brian): when offline instantiate the BlockService with a bitswap
@@ -256,41 +256,6 @@ func initIdentity(cfg *config.Identity, peers peer.Peerstore, online bool) (peer
256 return self, nil
257 }
258
259 -func initConnections(ctx context.Context, bootstrap []*config.BootstrapPeer, pstore peer.Peerstore, route *dht.IpfsDHT) {
260 - // TODO consider stricter error handling
261 - // TODO consider Criticalf error logging
262 - for _, p := range bootstrap {
263 - if p.PeerID == "" {
264 - log.Criticalf("error: peer does not include PeerID. %v", p)
265 - }
266 -
267 - maddr, err := ma.NewMultiaddr(p.Address)
268 - if err != nil {
269 - log.Error(err)
270 - continue
271 - }
272 -
273 - // setup peer
274 - id, err := peer.DecodePrettyID(p.PeerID)
275 - if err != nil {
276 - log.Criticalf("Bootstrapping error: %v", err)
277 - continue
278 - }
279 - npeer, err := pstore.FindOrCreate(id)
280 - if err != nil {
281 - log.Criticalf("Bootstrapping error: %v", err)
282 - continue
283 - }
284 - npeer.AddAddress(maddr)
285 -
286 - if err := route.Connect(ctx, npeer); err != nil {
287 - log.Criticalf("Bootstrapping error: %v", err)
288 - continue
289 - }
290 - log.Event(ctx, "bootstrap", npeer)
291 - }
292 -}
293 -
259 func listenAddresses(cfg *config.Config) ([]ma.Multiaddr, error) {
260
261 var err error