@cryptotaxi247 / kubo / commits / 36494e493

ipns(pubsub): utilize persistent pubsub value store

Adin Schmahmann committed Nov 11, 2019 at 11:17 UTC 36494e493a18441ecf268cbe6cf3c1e45b70a06c
6 files changed +54 -15
core/node/groups.go
+4 -1
@@ -67,8 +67,10 @@ func LibP2P(bcfg *BuildCfg, cfg *config.Config) fx.Option {
67
68 // parse PubSub config
69
70 - ps := fx.Options()
70 + ps, disc := fx.Options(), fx.Options()
71 if bcfg.getOpt("pubsub") || bcfg.getOpt("ipnsps") {
72 + disc = fx.Provide(libp2p.TopicDiscovery())
73 +
74 var pubsubOptions []pubsub.Option
75 pubsubOptions = append(
76 pubsubOptions,
@@ -113,6 +115,7 @@ func LibP2P(bcfg *BuildCfg, cfg *config.Config) fx.Option {
115 maybeInvoke(libp2p.AutoNATService(cfg.Experimental.QUIC), cfg.Swarm.EnableAutoNATService),
116 connmgr,
117 ps,
118 + disc,
119 )
120
121 return opts
core/node/libp2p/pubsub.go
+6 -5
@@ -1,7 +1,8 @@
1 package libp2p
2
3 import (
4 - host "github.com/libp2p/go-libp2p-core/host"
4 + "github.com/libp2p/go-libp2p-core/discovery"
5 + "github.com/libp2p/go-libp2p-core/host"
6 pubsub "github.com/libp2p/go-libp2p-pubsub"
7 "go.uber.org/fx"
8
@@ -9,13 +10,13 @@ import (
10 )
11
12 func FloodSub(pubsubOptions ...pubsub.Option) interface{} {
12 - return func(mctx helpers.MetricsCtx, lc fx.Lifecycle, host host.Host) (service *pubsub.PubSub, err error) {
13 - return pubsub.NewFloodSub(helpers.LifecycleCtx(mctx, lc), host, pubsubOptions...)
13 + return func(mctx helpers.MetricsCtx, lc fx.Lifecycle, host host.Host, disc discovery.Discovery) (service *pubsub.PubSub, err error) {
14 + return pubsub.NewFloodSub(helpers.LifecycleCtx(mctx, lc), host, append(pubsubOptions, pubsub.WithDiscovery(disc))...)
15 }
16 }
17
18 func GossipSub(pubsubOptions ...pubsub.Option) interface{} {
18 - return func(mctx helpers.MetricsCtx, lc fx.Lifecycle, host host.Host) (service *pubsub.PubSub, err error) {
19 - return pubsub.NewGossipSub(helpers.LifecycleCtx(mctx, lc), host, pubsubOptions...)
19 + return func(mctx helpers.MetricsCtx, lc fx.Lifecycle, host host.Host, disc discovery.Discovery) (service *pubsub.PubSub, err error) {
20 + return pubsub.NewGossipSub(helpers.LifecycleCtx(mctx, lc), host, append(pubsubOptions, pubsub.WithDiscovery(disc))...)
21 }
22 }
core/node/libp2p/routing.go
+9 -4
@@ -3,6 +3,7 @@ package libp2p
3 import (
4 "context"
5 "sort"
6 + "time"
7
8 host "github.com/libp2p/go-libp2p-core/host"
9 routing "github.com/libp2p/go-libp2p-core/routing"
@@ -83,15 +84,19 @@ type p2pPSRoutingIn struct {
84 PubSub *pubsub.PubSub `optional:"true"`
85 }
86
86 -func PubsubRouter(mctx helpers.MetricsCtx, lc fx.Lifecycle, in p2pPSRoutingIn) (p2pRouterOut, *namesys.PubsubValueStore) {
87 - psRouter := namesys.NewPubsubValueStore(
87 +func PubsubRouter(mctx helpers.MetricsCtx, lc fx.Lifecycle, in p2pPSRoutingIn) (p2pRouterOut, *namesys.PubsubValueStore, error) {
88 + psRouter, err := namesys.NewPubsubValueStore(
89 helpers.LifecycleCtx(mctx, lc),
90 in.Host,
90 - in.BaseIpfsRouting,
91 in.PubSub,
92 in.Validator,
93 + namesys.WithRebroadcastInterval(time.Minute),
94 )
95
96 + if err != nil {
97 + return p2pRouterOut{}, nil, err
98 + }
99 +
100 return p2pRouterOut{
101 Router: Router{
102 Routing: &routinghelpers.Compose{
@@ -102,5 +107,5 @@ func PubsubRouter(mctx helpers.MetricsCtx, lc fx.Lifecycle, in p2pPSRoutingIn) (
107 },
108 Priority: 100,
109 },
105 - }, psRouter
110 + }, psRouter, nil
111 }
core/node/libp2p/topicdiscovery.go new
+31
@@ -0,0 +1,31 @@
1 +package libp2p
2 +
3 +import (
4 + "math/rand"
5 + "time"
6 +
7 + "github.com/libp2p/go-libp2p-core/discovery"
8 + "github.com/libp2p/go-libp2p-core/host"
9 + disc "github.com/libp2p/go-libp2p-discovery"
10 +
11 + "github.com/ipfs/go-ipfs/core/node/helpers"
12 + "go.uber.org/fx"
13 +)
14 +
15 +func TopicDiscovery() interface{} {
16 + return func(mctx helpers.MetricsCtx, lc fx.Lifecycle, host host.Host, cr BaseIpfsRouting) (service discovery.Discovery, err error) {
17 + baseDisc := disc.NewRoutingDiscovery(cr)
18 + minBackoff, maxBackoff := time.Second*60, time.Hour
19 + rng := rand.New(rand.NewSource(rand.Int63()))
20 + d, err := disc.NewBackoffDiscovery(
21 + baseDisc,
22 + disc.NewExponentialBackoff(minBackoff, maxBackoff, disc.FullJitter, time.Second, 5.0, 0, rng),
23 + )
24 +
25 + if err != nil {
26 + return nil, err
27 + }
28 +
29 + return d, nil
30 + }
31 +}
go.mod
+2 -1
@@ -62,6 +62,7 @@ require (
62 github.com/libp2p/go-libp2p-circuit v0.1.4
63 github.com/libp2p/go-libp2p-connmgr v0.1.1
64 github.com/libp2p/go-libp2p-core v0.2.5
65 + github.com/libp2p/go-libp2p-discovery v0.2.0
66 github.com/libp2p/go-libp2p-http v0.1.4
67 github.com/libp2p/go-libp2p-kad-dht v0.3.1
68 github.com/libp2p/go-libp2p-kbucket v0.2.1
@@ -70,7 +71,7 @@ require (
71 github.com/libp2p/go-libp2p-peerstore v0.1.4
72 github.com/libp2p/go-libp2p-pnet v0.1.0
73 github.com/libp2p/go-libp2p-pubsub v0.2.4
73 - github.com/libp2p/go-libp2p-pubsub-router v0.1.0
74 + github.com/libp2p/go-libp2p-pubsub-router v0.2.0
75 github.com/libp2p/go-libp2p-quic-transport v0.2.2
76 github.com/libp2p/go-libp2p-record v0.1.2
77 github.com/libp2p/go-libp2p-routing-helpers v0.1.0
go.sum
+2 -4
@@ -168,7 +168,6 @@ github.com/ipfs/go-datastore v0.1.0 h1:TOxI04l8CmO4zGtesENhzm4PwkFwJXY3rKiYaaMf9
168 github.com/ipfs/go-datastore v0.1.0/go.mod h1:d4KVXhMt913cLBEI/PXAy6ko+W7e9AhyAKBGh803qeE=
169 github.com/ipfs/go-datastore v0.1.1 h1:F4k0TkTAZGLFzBOrVKDAvch6JZtuN4NHkfdcEZL50aI=
170 github.com/ipfs/go-datastore v0.1.1/go.mod h1:w38XXW9kVFNp57Zj5knbKWM2T+KOZCGDRVNdgPHtbHw=
171 -github.com/ipfs/go-datastore v0.3.0 h1:9au0tYi/+n7xeUnGHG6davnS8x9hWbOzP/388Vx3CMs=
171 github.com/ipfs/go-datastore v0.3.0/go.mod h1:w38XXW9kVFNp57Zj5knbKWM2T+KOZCGDRVNdgPHtbHw=
172 github.com/ipfs/go-datastore v0.3.1 h1:SS1t869a6cctoSYmZXUk8eL6AzVXgASmKIWFNQkQ1jU=
173 github.com/ipfs/go-datastore v0.3.1/go.mod h1:w38XXW9kVFNp57Zj5knbKWM2T+KOZCGDRVNdgPHtbHw=
@@ -433,11 +432,10 @@ github.com/libp2p/go-libp2p-pnet v0.1.0 h1:kRUES28dktfnHNIRW4Ro78F7rKBHBiw5MJpl0
432 github.com/libp2p/go-libp2p-pnet v0.1.0/go.mod h1:ZkyZw3d0ZFOex71halXRihWf9WH/j3OevcJdTmD0lyE=
433 github.com/libp2p/go-libp2p-protocol v0.0.1/go.mod h1:Af9n4PiruirSDjHycM1QuiMi/1VZNHYcK8cLgFJLZ4s=
434 github.com/libp2p/go-libp2p-protocol v0.1.0/go.mod h1:KQPHpAabB57XQxGrXCNvbL6UEXfQqUgC/1adR2Xtflk=
436 -github.com/libp2p/go-libp2p-pubsub v0.1.0/go.mod h1:ZwlKzRSe1eGvSIdU5bD7+8RZN/Uzw0t1Bp9R1znpR/Q=
435 github.com/libp2p/go-libp2p-pubsub v0.2.4 h1:O4BcaKpPQ9p82yTBtzIzgDFoOXkqhrQpfcVac3FAywU=
436 github.com/libp2p/go-libp2p-pubsub v0.2.4/go.mod h1:1tJwAfySvZQ49R9uTVlkwtSTMVLeQQdrnLTJrr91gVc=
439 -github.com/libp2p/go-libp2p-pubsub-router v0.1.0 h1:xA5B8Sdx64tNlSRIcay2QUngtlu8LpUJClaUk/dYYrg=
440 -github.com/libp2p/go-libp2p-pubsub-router v0.1.0/go.mod h1:PnHOshBr/2I2ZxVfEsqfgCQPsVg09zo+DhSlWkOhPFM=
437 +github.com/libp2p/go-libp2p-pubsub-router v0.2.0 h1:AcrCIL2aUNiQkQuqKonqVxhFBEfP0PsHVUdAowZZVLA=
438 +github.com/libp2p/go-libp2p-pubsub-router v0.2.0/go.mod h1:CgbGriQhei3gy9y3MwmAapge+SYGYNrcwOeXp3Sefpg=
439 github.com/libp2p/go-libp2p-quic-transport v0.2.2 h1:XyGRqFHD1oHdI2k98P1tWWRb9s27fl1SfmCcaX8plso=
440 github.com/libp2p/go-libp2p-quic-transport v0.2.2/go.mod h1:rVzcsiuOFBomAqvNOxeBUcP4vM4wE+NqqRZWvxjkbe0=
441 github.com/libp2p/go-libp2p-record v0.0.1/go.mod h1:grzqg263Rug/sRex85QrDOLntdFAymLDLm7lxMgU79Q=