@cryptotaxi247 / kubo / commits / a7ef09554

refactor(bitswap:notifications) move, rename

add interface

Brian Tiger Chow committed Sep 12, 2014 at 03:08 UTC a7ef09554f92e43e8875337b8040d667b1658a3e
3 files changed +20 -13
bitswap/bitswap.go
+3 -2
@@ -7,6 +7,7 @@ import (
7 proto "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/goprotobuf/proto"
8 ds "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/datastore.go"
9
10 + notifications "github.com/jbenet/go-ipfs/bitswap/notifications"
11 blocks "github.com/jbenet/go-ipfs/blocks"
12 swarm "github.com/jbenet/go-ipfs/net/swarm"
13 peer "github.com/jbenet/go-ipfs/peer"
@@ -39,7 +40,7 @@ type BitSwap struct {
40 // routing interface for communication
41 routing *dht.IpfsDHT
42
42 - notifications *notifications
43 + notifications notifications.PubSub
44
45 // partners is a map of currently active bitswap relationships.
46 // The Ledger has the peer.ID, and the peer connection works through net.
@@ -69,7 +70,7 @@ func NewBitSwap(p *peer.Peer, net swarm.Network, d ds.Datastore, r routing.IpfsR
70 routing: r.(*dht.IpfsDHT),
71 meschan: net.GetChannel(swarm.PBWrapper_BITSWAP),
72 haltChan: make(chan struct{}),
72 - notifications: newNotifications(),
73 + notifications: notifications.New(),
74 }
75
76 go bs.handleMessages()
bitswap/notifications/notifications.go renamed
+14 -8
@@ -1,4 +1,4 @@
1 -package bitswap
1 +package notifications
2
3 import (
4 context "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
@@ -8,16 +8,22 @@ import (
8 u "github.com/jbenet/go-ipfs/util"
9 )
10
11 -type notifications struct {
12 - wrapped *pubsub.PubSub
11 +type PubSub interface {
12 + Publish(block *blocks.Block)
13 + Subscribe(ctx context.Context, k u.Key) <-chan *blocks.Block
14 + Shutdown()
15 }
16
15 -func newNotifications() *notifications {
17 +func New() PubSub {
18 const bufferSize = 16
17 - return &notifications{pubsub.New(bufferSize)}
19 + return &impl{pubsub.New(bufferSize)}
20 +}
21 +
22 +type impl struct {
23 + wrapped *pubsub.PubSub
24 }
25
20 -func (ps *notifications) Publish(block *blocks.Block) {
26 +func (ps *impl) Publish(block *blocks.Block) {
27 topic := string(block.Key())
28 ps.wrapped.Pub(block, topic)
29 }
@@ -25,7 +31,7 @@ func (ps *notifications) Publish(block *blocks.Block) {
31 // Subscribe returns a one-time use |blockChannel|. |blockChannel| returns nil
32 // if the |ctx| times out or is cancelled. Then channel is closed after the
33 // block given by |k| is sent.
28 -func (ps *notifications) Subscribe(ctx context.Context, k u.Key) <-chan *blocks.Block {
34 +func (ps *impl) Subscribe(ctx context.Context, k u.Key) <-chan *blocks.Block {
35 topic := string(k)
36 subChan := ps.wrapped.SubOnce(topic)
37 blockChannel := make(chan *blocks.Block)
@@ -44,6 +50,6 @@ func (ps *notifications) Subscribe(ctx context.Context, k u.Key) <-chan *blocks.
50 return blockChannel
51 }
52
47 -func (ps *notifications) Shutdown() {
53 +func (ps *impl) Shutdown() {
54 ps.wrapped.Shutdown()
55 }
bitswap/notifications/notifications_test.go renamed
+3 -3
@@ -1,4 +1,4 @@
1 -package bitswap
1 +package notifications
2
3 import (
4 "bytes"
@@ -13,7 +13,7 @@ import (
13 func TestPublishSubscribe(t *testing.T) {
14 blockSent := getBlockOrFail(t, "Greetings from The Interval")
15
16 - n := newNotifications()
16 + n := New()
17 defer n.Shutdown()
18 ch := n.Subscribe(context.Background(), blockSent.Key())
19
@@ -28,7 +28,7 @@ func TestCarryOnWhenDeadlineExpires(t *testing.T) {
28 impossibleDeadline := time.Nanosecond
29 fastExpiringCtx, _ := context.WithTimeout(context.Background(), impossibleDeadline)
30
31 - n := newNotifications()
31 + n := New()
32 defer n.Shutdown()
33 blockChannel := n.Subscribe(fastExpiringCtx, getBlockOrFail(t, "A Missed Connection").Key())
34