@cryptotaxi247 / kubo / commits / ff770fadc

add discovery option, and update to floodsub 0.4.2

License: MIT Signed-off-by: Jeromy <why@ipfs.io>

Jeromy committed Sep 13, 2016 at 10:19 UTC ff770fadc896218f9fc9df3dcdf382669ad72c09
3 files changed +51 -4
core/commands/pubsub.go
+48 -1
@@ -2,13 +2,21 @@ package commands
2
3 import (
4 "bytes"
5 + "context"
6 "encoding/binary"
7 "io"
8 + "sync"
9 + "time"
10
11 + blocks "github.com/ipfs/go-ipfs/blocks"
12 cmds "github.com/ipfs/go-ipfs/commands"
13 + core "github.com/ipfs/go-ipfs/core"
14
10 - floodsub "gx/ipfs/QmSWp1Yx7Z5pbpeCbUy6tfFj2DrHUe7tGQqyYC2vspbXH1/floodsub"
15 + floodsub "gx/ipfs/QmQtsU1T46uxjFMd5r5PfyaY1HdV5jcxZbvvHbAVRL52hc/floodsub"
16 u "gx/ipfs/QmZNVWh8LLjAavuQ2JXuFmuYH3C11xo988vSgp7UQrTRj1/go-ipfs-util"
17 + key "gx/ipfs/Qmce4Y4zg3sYr7xKM5UueS67vhNni6EeWgCRnb7MbLJMew/go-key"
18 + pstore "gx/ipfs/QmdMfSLMDBDYhtc4oF3NYGCZr5dy4wQb6Ji26N4D4mdxa2/go-libp2p-peerstore"
19 + cid "gx/ipfs/QmfSc2xehWmWLnwwYR91Y8QF4xdASypTFVknutoKQS3GHp/go-cid"
20 )
21
22 var PubsubCmd = &cmds.Command{
@@ -41,6 +49,9 @@ to be used in a production environment.
49 Arguments: []cmds.Argument{
50 cmds.StringArg("topic", true, false, "String name of topic to subscribe to."),
51 },
52 + Options: []cmds.Option{
53 + cmds.BoolOption("discover", "try to discover other peers subscribed to the same topic"),
54 + },
55 Run: func(req cmds.Request, res cmds.Response) {
56 n, err := req.InvocContext().GetNode()
57 if err != nil {
@@ -79,6 +90,18 @@ to be used in a production environment.
90 }
91 }
92 }()
93 +
94 + discover, _, _ := req.Option("discover").Bool()
95 + if discover {
96 + blk := blocks.NewBlock([]byte("floodsub:" + topic))
97 + cid, err := n.Blocks.AddObject(blk)
98 + if err != nil {
99 + log.Error("pubsub discovery: ", err)
100 + return
101 + }
102 +
103 + connectToPubSubPeers(req.Context(), n, cid)
104 + }
105 },
106 Marshalers: cmds.MarshalerMap{
107 cmds.Text: getPsMsgMarshaler(func(m *floodsub.Message) (io.Reader, error) {
@@ -97,6 +120,30 @@ to be used in a production environment.
120 Type: floodsub.Message{},
121 }
122
123 +func connectToPubSubPeers(ctx context.Context, n *core.IpfsNode, cid *cid.Cid) {
124 + ctx, cancel := context.WithCancel(ctx)
125 + defer cancel()
126 +
127 + provs := n.Routing.FindProvidersAsync(ctx, key.Key(cid.Hash()), 10)
128 + wg := &sync.WaitGroup{}
129 + for p := range provs {
130 + wg.Add(1)
131 + go func(pi pstore.PeerInfo) {
132 + defer wg.Done()
133 + ctx, cancel := context.WithTimeout(ctx, time.Second*10)
134 + defer cancel()
135 + err := n.PeerHost.Connect(ctx, pi)
136 + if err != nil {
137 + log.Info("pubsub discover: ", err)
138 + return
139 + }
140 + log.Info("connected to pubsub peer:", pi.ID)
141 + }(p)
142 + }
143 +
144 + wg.Wait()
145 +}
146 +
147 func getPsMsgMarshaler(f func(m *floodsub.Message) (io.Reader, error)) func(cmds.Response) (io.Reader, error) {
148 return func(res cmds.Response) (io.Reader, error) {
149 outChan, ok := res.Output().(<-chan interface{})
core/core.go
+1 -1
@@ -17,9 +17,9 @@ import (
17 "time"
18
19 diag "github.com/ipfs/go-ipfs/diagnostics"
20 + floodsub "gx/ipfs/QmQtsU1T46uxjFMd5r5PfyaY1HdV5jcxZbvvHbAVRL52hc/floodsub"
21 goprocess "gx/ipfs/QmSF8fPo3jgVBAy8fpdjjYqgG87dkJgUprRBHRd2tmfgpP/goprocess"
22 mamask "gx/ipfs/QmSMZwvs3n4GBikZ7hKzT17c3bk65FmyZo2JqtJ16swqCv/multiaddr-filter"
22 - floodsub "gx/ipfs/QmSWp1Yx7Z5pbpeCbUy6tfFj2DrHUe7tGQqyYC2vspbXH1/floodsub"
23 logging "gx/ipfs/QmSpJByNKFX1sCsHBEp3R73FL4NF6FnQTEGyNAXHm2GS52/go-log"
24 b58 "gx/ipfs/QmT8rehPR3F6bmwL6zjUN8XpiDBFFpMP2myPdC6ApsWfJf/go-base58"
25 ic "gx/ipfs/QmVoi5es8D5fNHZDqoW6DgDAEPEV5hQp8GBz161vZXiwpQ/go-libp2p-crypto"
package.json
+2 -2
@@ -266,9 +266,9 @@
266 },
267 {
268 "author": "whyrusleeping",
269 - "hash": "QmSWp1Yx7Z5pbpeCbUy6tfFj2DrHUe7tGQqyYC2vspbXH1",
269 + "hash": "QmQtsU1T46uxjFMd5r5PfyaY1HdV5jcxZbvvHbAVRL52hc",
270 "name": "floodsub",
271 - "version": "0.4.1"
271 + "version": "0.4.2"
272 }
273 ],
274 "gxVersion": "0.4.0",