use floodsub version 0.8.1
License: MIT Signed-off-by: Jan Winkelmann <j-winkelmann@tuhh.de>
Jan Winkelmann committed
Nov 21, 2016 at 16:44 UTC
05265f176bc586f126a4a2c50548933ccf46903c
3 files changed
+14
-14
core/commands/pubsub.go
+11
-11
@@ -13,7 +13,7 @@ import (
13
cmds "github.com/ipfs/go-ipfs/commands"
14
core "github.com/ipfs/go-ipfs/core"
15
16
- floodsub "gx/ipfs/QmWiLbk7eE1jGePDAuS26E2A9bMK3e3PMH3dcSeRY3MEBR/floodsub"
16
+ floodsub "gx/ipfs/QmRJs5veT3gnuYpLAagC3NbzixbkgwjSdUXTKfh3hMo6XM/floodsub"
17
pstore "gx/ipfs/QmXXCcQ7CLg5a81Ui9TTR35QcR4y7ZyihxwfjqaHfUVcVo/go-libp2p-peerstore"
18
u "gx/ipfs/Qmb912gdngC1UWwTkhuW8knyRbcWeu5kqkxBpveLmW8bSr/go-ipfs-util"
19
cid "gx/ipfs/QmcEcrBAMrwMyhSjXt4yfyPpzgSuV8HLHavnfmiKCSRqZU/go-cid"
@@ -77,7 +77,7 @@ To use, the daemon must be run with '--enable-pubsub-experiment'.
77
}
78
79
topic := req.Arguments()[0]
80
- msgs, err := n.Floodsub.Subscribe(req.Context(), topic)
80
+ sub, err := n.Floodsub.Subscribe(topic)
81
if err != nil {
82
res.SetError(err, cmds.ErrNormal)
83
return
@@ -86,19 +86,19 @@ To use, the daemon must be run with '--enable-pubsub-experiment'.
86
out := make(chan interface{})
87
res.SetOutput((<-chan interface{})(out))
88
89
- ctx := req.Context()
89
go func() {
90
+ defer sub.Cancel()
91
defer close(out)
92
for {
93
- select {
94
- case msg, ok := <-msgs:
95
- if !ok {
96
- return
97
- }
98
- out <- msg
99
- case <-ctx.Done():
100
- n.Floodsub.Unsub(topic)
93
+ msg, err := sub.Next(req.Context())
94
+ if err == io.EOF || err == context.Canceled {
95
+ break
96
+ } else if err != nil {
97
+ res.SetError(err, cmds.ErrNormal)
98
+ return
99
}
100
+
101
+ out <- msg
102
}
103
}()
104
core/core.go
+1
-1
@@ -38,6 +38,7 @@ import (
38
ft "github.com/ipfs/go-ipfs/unixfs"
39
40
swarm "gx/ipfs/QmNafAGBU21iQmLudMT2z1kqgEGhjUrNoK9a3v4azd8ei4/go-libp2p-swarm"
41
+ floodsub "gx/ipfs/QmRJs5veT3gnuYpLAagC3NbzixbkgwjSdUXTKfh3hMo6XM/floodsub"
42
goprocess "gx/ipfs/QmSF8fPo3jgVBAy8fpdjjYqgG87dkJgUprRBHRd2tmfgpP/goprocess"
43
mamask "gx/ipfs/QmSMZwvs3n4GBikZ7hKzT17c3bk65FmyZo2JqtJ16swqCv/multiaddr-filter"
44
logging "gx/ipfs/QmSpJByNKFX1sCsHBEp3R73FL4NF6FnQTEGyNAXHm2GS52/go-log"
@@ -45,7 +46,6 @@ import (
46
ma "gx/ipfs/QmUAQaWbKxGCUTuoQVvvicbQNZ9APF5pDGWyAZSe93AtKH/go-multiaddr"
47
routing "gx/ipfs/QmUrCwTDvJgmBbJVHu1HGEyqDaod3dR6sEkZkpxZk4u47c/go-libp2p-routing"
48
addrutil "gx/ipfs/QmVDnc2zvyQm8LhT72n22THcshvH7j3qPMnhvjerQER62T/go-addr-util"
48
- floodsub "gx/ipfs/QmWiLbk7eE1jGePDAuS26E2A9bMK3e3PMH3dcSeRY3MEBR/floodsub"
49
metrics "gx/ipfs/QmX4j1JhubdEt4EB1JY1mMKTvJwPZSRzTv3uwh5zaDqyAi/go-libp2p-metrics"
50
pstore "gx/ipfs/QmXXCcQ7CLg5a81Ui9TTR35QcR4y7ZyihxwfjqaHfUVcVo/go-libp2p-peerstore"
51
discovery "gx/ipfs/QmZyBJGpRnbQ7oUstoGNZbhXC4HJuFUCgpp8pmsVTUwdS3/go-libp2p/p2p/discovery"
package.json
+2
-2
@@ -266,9 +266,9 @@
266
},
267
{
268
"author": "whyrusleeping",
269
- "hash": "QmWiLbk7eE1jGePDAuS26E2A9bMK3e3PMH3dcSeRY3MEBR",
269
+ "hash": "QmRJs5veT3gnuYpLAagC3NbzixbkgwjSdUXTKfh3hMo6XM",
270
"name": "floodsub",
271
- "version": "0.8.0"
271
+ "version": "0.8.1"
272
},
273
{
274
"author": "whyrusleeping",