@cryptotaxi247 / kubo / commits / 455c6582f

refactor(bitswap) leverage third-party pubsub lib

use a third-party pubsub library for internal communications Insights: * Within bitswap, the actors don't need anything more than simple pubsub behavior. Wrapping and unwrapping messages proves unneccessary. Changes: * Simplifies the interface for both actors calling GetBlock and actors receiving blocks on the network * Leverages a well-tested third-party pubsub library Design Goals: * reduce complexity * extract implementation details (wrapping and unwrapping data, etc) from bitswap and let bitswap focus on composition of core algorithms operations

Brian Tiger Chow committed Sep 11, 2014 at 16:15 UTC 455c6582f52fc14bfec1a9051b7a8592303d9684
6 files changed +582
Godeps/Godeps.json
+4
@@ -71,6 +71,10 @@
71 {
72 "ImportPath": "github.com/syndtr/goleveldb/leveldb",
73 "Rev": "99056d50e56252fbe0021d5c893defca5a76baf8"
74 + },
75 + {
76 + "ImportPath": "github.com/tuxychandru/pubsub",
77 + "Rev": "02de8aa2db3d570c5ab1be5ba67b456fd0fb7c4e"
78 }
79 ]
80 }
Godeps/_workspace/src/github.com/tuxychandru/pubsub/README.md new
+30
@@ -0,0 +1,30 @@
1 +Install pubsub with,
2 +
3 + go get github.com/tuxychandru/pubsub
4 +
5 +View the [API Documentation](http://godoc.org/github.com/tuxychandru/pubsub).
6 +
7 +## License
8 +
9 +Copyright (c) 2013, Chandra Sekar S
10 +All rights reserved.
11 +
12 +Redistribution and use in source and binary forms, with or without
13 +modification, are permitted provided that the following conditions are met:
14 +
15 +1. Redistributions of source code must retain the above copyright notice, this
16 + list of conditions and the following disclaimer.
17 +2. Redistributions in binary form must reproduce the above copyright notice,
18 + this list of conditions and the following disclaimer in the documentation
19 + and/or other materials provided with the distribution.
20 +
21 +THIS SOFTWARE IS PROVIDED BY THE COPYRIGHT HOLDERS AND CONTRIBUTORS "AS IS" AND
22 +ANY EXPRESS OR IMPLIED WARRANTIES, INCLUDING, BUT NOT LIMITED TO, THE IMPLIED
23 +WARRANTIES OF MERCHANTABILITY AND FITNESS FOR A PARTICULAR PURPOSE ARE
24 +DISCLAIMED. IN NO EVENT SHALL THE COPYRIGHT HOLDER OR CONTRIBUTORS BE LIABLE FOR
25 +ANY DIRECT, INDIRECT, INCIDENTAL, SPECIAL, EXEMPLARY, OR CONSEQUENTIAL DAMAGES
26 +(INCLUDING, BUT NOT LIMITED TO, PROCUREMENT OF SUBSTITUTE GOODS OR SERVICES;
27 +LOSS OF USE, DATA, OR PROFITS; OR BUSINESS INTERRUPTION) HOWEVER CAUSED AND
28 +ON ANY THEORY OF LIABILITY, WHETHER IN CONTRACT, STRICT LIABILITY, OR TORT
29 +(INCLUDING NEGLIGENCE OR OTHERWISE) ARISING IN ANY WAY OUT OF THE USE OF THIS
30 +SOFTWARE, EVEN IF ADVISED OF THE POSSIBILITY OF SUCH DAMAGE.
Godeps/_workspace/src/github.com/tuxychandru/pubsub/pubsub.go new
+208
@@ -0,0 +1,208 @@
1 +// Copyright 2013, Chandra Sekar S. All rights reserved.
2 +// Use of this source code is governed by a BSD-style
3 +// license that can be found in the README.md file.
4 +
5 +// Package pubsub implements a simple multi-topic pub-sub
6 +// library.
7 +//
8 +// Topics must be strings and messages of any type can be
9 +// published. A topic can have any number of subcribers and
10 +// all of them receive messages published on the topic.
11 +package pubsub
12 +
13 +type operation int
14 +
15 +const (
16 + sub operation = iota
17 + subOnce
18 + pub
19 + unsub
20 + unsubAll
21 + closeTopic
22 + shutdown
23 +)
24 +
25 +// PubSub is a collection of topics.
26 +type PubSub struct {
27 + cmdChan chan cmd
28 + capacity int
29 +}
30 +
31 +type cmd struct {
32 + op operation
33 + topics []string
34 + ch chan interface{}
35 + msg interface{}
36 +}
37 +
38 +// New creates a new PubSub and starts a goroutine for handling operations.
39 +// The capacity of the channels created by Sub and SubOnce will be as specified.
40 +func New(capacity int) *PubSub {
41 + ps := &PubSub{make(chan cmd), capacity}
42 + go ps.start()
43 + return ps
44 +}
45 +
46 +// Sub returns a channel on which messages published on any of
47 +// the specified topics can be received.
48 +func (ps *PubSub) Sub(topics ...string) chan interface{} {
49 + return ps.sub(sub, topics...)
50 +}
51 +
52 +// SubOnce is similar to Sub, but only the first message published, after subscription,
53 +// on any of the specified topics can be received.
54 +func (ps *PubSub) SubOnce(topics ...string) chan interface{} {
55 + return ps.sub(subOnce, topics...)
56 +}
57 +
58 +func (ps *PubSub) sub(op operation, topics ...string) chan interface{} {
59 + ch := make(chan interface{}, ps.capacity)
60 + ps.cmdChan <- cmd{op: op, topics: topics, ch: ch}
61 + return ch
62 +}
63 +
64 +// AddSub adds subscriptions to an existing channel.
65 +func (ps *PubSub) AddSub(ch chan interface{}, topics ...string) {
66 + ps.cmdChan <- cmd{op: sub, topics: topics, ch: ch}
67 +}
68 +
69 +// Pub publishes the given message to all subscribers of
70 +// the specified topics.
71 +func (ps *PubSub) Pub(msg interface{}, topics ...string) {
72 + ps.cmdChan <- cmd{op: pub, topics: topics, msg: msg}
73 +}
74 +
75 +// Unsub unsubscribes the given channel from the specified
76 +// topics. If no topic is specified, it is unsubscribed
77 +// from all topics.
78 +func (ps *PubSub) Unsub(ch chan interface{}, topics ...string) {
79 + if len(topics) == 0 {
80 + ps.cmdChan <- cmd{op: unsubAll, ch: ch}
81 + return
82 + }
83 +
84 + ps.cmdChan <- cmd{op: unsub, topics: topics, ch: ch}
85 +}
86 +
87 +// Close closes all channels currently subscribed to the specified topics.
88 +// If a channel is subscribed to multiple topics, some of which is
89 +// not specified, it is not closed.
90 +func (ps *PubSub) Close(topics ...string) {
91 + ps.cmdChan <- cmd{op: closeTopic, topics: topics}
92 +}
93 +
94 +// Shutdown closes all subscribed channels and terminates the goroutine.
95 +func (ps *PubSub) Shutdown() {
96 + ps.cmdChan <- cmd{op: shutdown}
97 +}
98 +
99 +func (ps *PubSub) start() {
100 + reg := registry{
101 + topics: make(map[string]map[chan interface{}]bool),
102 + revTopics: make(map[chan interface{}]map[string]bool),
103 + }
104 +
105 +loop:
106 + for cmd := range ps.cmdChan {
107 + if cmd.topics == nil {
108 + switch cmd.op {
109 + case unsubAll:
110 + reg.removeChannel(cmd.ch)
111 +
112 + case shutdown:
113 + break loop
114 + }
115 +
116 + continue loop
117 + }
118 +
119 + for _, topic := range cmd.topics {
120 + switch cmd.op {
121 + case sub:
122 + reg.add(topic, cmd.ch, false)
123 +
124 + case subOnce:
125 + reg.add(topic, cmd.ch, true)
126 +
127 + case pub:
128 + reg.send(topic, cmd.msg)
129 +
130 + case unsub:
131 + reg.remove(topic, cmd.ch)
132 +
133 + case closeTopic:
134 + reg.removeTopic(topic)
135 + }
136 + }
137 + }
138 +
139 + for topic, chans := range reg.topics {
140 + for ch, _ := range chans {
141 + reg.remove(topic, ch)
142 + }
143 + }
144 +}
145 +
146 +// registry maintains the current subscription state. It's not
147 +// safe to access a registry from multiple goroutines simultaneously.
148 +type registry struct {
149 + topics map[string]map[chan interface{}]bool
150 + revTopics map[chan interface{}]map[string]bool
151 +}
152 +
153 +func (reg *registry) add(topic string, ch chan interface{}, once bool) {
154 + if reg.topics[topic] == nil {
155 + reg.topics[topic] = make(map[chan interface{}]bool)
156 + }
157 + reg.topics[topic][ch] = once
158 +
159 + if reg.revTopics[ch] == nil {
160 + reg.revTopics[ch] = make(map[string]bool)
161 + }
162 + reg.revTopics[ch][topic] = true
163 +}
164 +
165 +func (reg *registry) send(topic string, msg interface{}) {
166 + for ch, once := range reg.topics[topic] {
167 + ch <- msg
168 + if once {
169 + for topic := range reg.revTopics[ch] {
170 + reg.remove(topic, ch)
171 + }
172 + }
173 + }
174 +}
175 +
176 +func (reg *registry) removeTopic(topic string) {
177 + for ch := range reg.topics[topic] {
178 + reg.remove(topic, ch)
179 + }
180 +}
181 +
182 +func (reg *registry) removeChannel(ch chan interface{}) {
183 + for topic := range reg.revTopics[ch] {
184 + reg.remove(topic, ch)
185 + }
186 +}
187 +
188 +func (reg *registry) remove(topic string, ch chan interface{}) {
189 + if _, ok := reg.topics[topic]; !ok {
190 + return
191 + }
192 +
193 + if _, ok := reg.topics[topic][ch]; !ok {
194 + return
195 + }
196 +
197 + delete(reg.topics[topic], ch)
198 + delete(reg.revTopics[ch], topic)
199 +
200 + if len(reg.topics[topic]) == 0 {
201 + delete(reg.topics, topic)
202 + }
203 +
204 + if len(reg.revTopics[ch]) == 0 {
205 + close(ch)
206 + delete(reg.revTopics, ch)
207 + }
208 +}
Godeps/_workspace/src/github.com/tuxychandru/pubsub/pubsub_test.go new
+230
@@ -0,0 +1,230 @@
1 +// Copyright 2013, Chandra Sekar S. All rights reserved.
2 +// Use of this source code is governed by a BSD-style
3 +// license that can be found in the README.md file.
4 +
5 +package pubsub
6 +
7 +import (
8 + check "launchpad.net/gocheck"
9 + "runtime"
10 + "testing"
11 + "time"
12 +)
13 +
14 +var _ = check.Suite(new(Suite))
15 +
16 +func Test(t *testing.T) {
17 + check.TestingT(t)
18 +}
19 +
20 +type Suite struct{}
21 +
22 +func (s *Suite) TestSub(c *check.C) {
23 + ps := New(1)
24 + ch1 := ps.Sub("t1")
25 + ch2 := ps.Sub("t1")
26 + ch3 := ps.Sub("t2")
27 +
28 + ps.Pub("hi", "t1")
29 + c.Check(<-ch1, check.Equals, "hi")
30 + c.Check(<-ch2, check.Equals, "hi")
31 +
32 + ps.Pub("hello", "t2")
33 + c.Check(<-ch3, check.Equals, "hello")
34 +
35 + ps.Shutdown()
36 + _, ok := <-ch1
37 + c.Check(ok, check.Equals, false)
38 + _, ok = <-ch2
39 + c.Check(ok, check.Equals, false)
40 + _, ok = <-ch3
41 + c.Check(ok, check.Equals, false)
42 +}
43 +
44 +func (s *Suite) TestSubOnce(c *check.C) {
45 + ps := New(1)
46 + ch := ps.SubOnce("t1")
47 +
48 + ps.Pub("hi", "t1")
49 + c.Check(<-ch, check.Equals, "hi")
50 +
51 + _, ok := <-ch
52 + c.Check(ok, check.Equals, false)
53 + ps.Shutdown()
54 +}
55 +
56 +func (s *Suite) TestAddSub(c *check.C) {
57 + ps := New(1)
58 + ch1 := ps.Sub("t1")
59 + ch2 := ps.Sub("t2")
60 +
61 + ps.Pub("hi1", "t1")
62 + c.Check(<-ch1, check.Equals, "hi1")
63 +
64 + ps.Pub("hi2", "t2")
65 + c.Check(<-ch2, check.Equals, "hi2")
66 +
67 + ps.AddSub(ch1, "t2", "t3")
68 + ps.Pub("hi3", "t2")
69 + c.Check(<-ch1, check.Equals, "hi3")
70 + c.Check(<-ch2, check.Equals, "hi3")
71 +
72 + ps.Pub("hi4", "t3")
73 + c.Check(<-ch1, check.Equals, "hi4")
74 +
75 + ps.Shutdown()
76 +}
77 +
78 +func (s *Suite) TestUnsub(c *check.C) {
79 + ps := New(1)
80 + ch := ps.Sub("t1")
81 +
82 + ps.Pub("hi", "t1")
83 + c.Check(<-ch, check.Equals, "hi")
84 +
85 + ps.Unsub(ch, "t1")
86 + _, ok := <-ch
87 + c.Check(ok, check.Equals, false)
88 + ps.Shutdown()
89 +}
90 +
91 +func (s *Suite) TestUnsubAll(c *check.C) {
92 + ps := New(1)
93 + ch1 := ps.Sub("t1", "t2", "t3")
94 + ch2 := ps.Sub("t1", "t3")
95 +
96 + ps.Unsub(ch1)
97 +
98 + m, ok := <-ch1
99 + c.Check(ok, check.Equals, false)
100 +
101 + ps.Pub("hi", "t1")
102 + m, ok = <-ch2
103 + c.Check(m, check.Equals, "hi")
104 +
105 + ps.Shutdown()
106 +}
107 +
108 +func (s *Suite) TestClose(c *check.C) {
109 + ps := New(1)
110 + ch1 := ps.Sub("t1")
111 + ch2 := ps.Sub("t1")
112 + ch3 := ps.Sub("t2")
113 + ch4 := ps.Sub("t3")
114 +
115 + ps.Pub("hi", "t1")
116 + ps.Pub("hello", "t2")
117 + c.Check(<-ch1, check.Equals, "hi")
118 + c.Check(<-ch2, check.Equals, "hi")
119 + c.Check(<-ch3, check.Equals, "hello")
120 +
121 + ps.Close("t1", "t2")
122 + _, ok := <-ch1
123 + c.Check(ok, check.Equals, false)
124 + _, ok = <-ch2
125 + c.Check(ok, check.Equals, false)
126 + _, ok = <-ch3
127 + c.Check(ok, check.Equals, false)
128 +
129 + ps.Pub("welcome", "t3")
130 + c.Check(<-ch4, check.Equals, "welcome")
131 +
132 + ps.Shutdown()
133 +}
134 +
135 +func (s *Suite) TestUnsubAfterClose(c *check.C) {
136 + ps := New(1)
137 + ch := ps.Sub("t1")
138 + defer func() {
139 + ps.Unsub(ch, "t1")
140 + ps.Shutdown()
141 + }()
142 +
143 + ps.Close("t1")
144 + _, ok := <-ch
145 + c.Check(ok, check.Equals, false)
146 +}
147 +
148 +func (s *Suite) TestShutdown(c *check.C) {
149 + start := runtime.NumGoroutine()
150 + New(10).Shutdown()
151 + time.Sleep(1)
152 + c.Check(runtime.NumGoroutine()-start, check.Equals, 1)
153 +}
154 +
155 +func (s *Suite) TestMultiSub(c *check.C) {
156 + ps := New(1)
157 + ch := ps.Sub("t1", "t2")
158 +
159 + ps.Pub("hi", "t1")
160 + c.Check(<-ch, check.Equals, "hi")
161 +
162 + ps.Pub("hello", "t2")
163 + c.Check(<-ch, check.Equals, "hello")
164 +
165 + ps.Shutdown()
166 + _, ok := <-ch
167 + c.Check(ok, check.Equals, false)
168 +}
169 +
170 +func (s *Suite) TestMultiSubOnce(c *check.C) {
171 + ps := New(1)
172 + ch := ps.SubOnce("t1", "t2")
173 +
174 + ps.Pub("hi", "t1")
175 + c.Check(<-ch, check.Equals, "hi")
176 +
177 + ps.Pub("hello", "t2")
178 +
179 + _, ok := <-ch
180 + c.Check(ok, check.Equals, false)
181 + ps.Shutdown()
182 +}
183 +
184 +func (s *Suite) TestMultiPub(c *check.C) {
185 + ps := New(1)
186 + ch1 := ps.Sub("t1")
187 + ch2 := ps.Sub("t2")
188 +
189 + ps.Pub("hi", "t1", "t2")
190 + c.Check(<-ch1, check.Equals, "hi")
191 + c.Check(<-ch2, check.Equals, "hi")
192 +
193 + ps.Shutdown()
194 +}
195 +
196 +func (s *Suite) TestMultiUnsub(c *check.C) {
197 + ps := New(1)
198 + ch := ps.Sub("t1", "t2", "t3")
199 +
200 + ps.Unsub(ch, "t1")
201 +
202 + ps.Pub("hi", "t1")
203 +
204 + ps.Pub("hello", "t2")
205 + c.Check(<-ch, check.Equals, "hello")
206 +
207 + ps.Unsub(ch, "t2", "t3")
208 + _, ok := <-ch
209 + c.Check(ok, check.Equals, false)
210 +
211 + ps.Shutdown()
212 +}
213 +
214 +func (s *Suite) TestMultiClose(c *check.C) {
215 + ps := New(1)
216 + ch := ps.Sub("t1", "t2")
217 +
218 + ps.Pub("hi", "t1")
219 + c.Check(<-ch, check.Equals, "hi")
220 +
221 + ps.Close("t1")
222 + ps.Pub("hello", "t2")
223 + c.Check(<-ch, check.Equals, "hello")
224 +
225 + ps.Close("t2")
226 + _, ok := <-ch
227 + c.Check(ok, check.Equals, false)
228 +
229 + ps.Shutdown()
230 +}
bitswap/notifications.go new
+50
@@ -0,0 +1,50 @@
1 +package bitswap
2 +
3 +import (
4 + context "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
5 + pubsub "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/tuxychandru/pubsub"
6 +
7 + blocks "github.com/jbenet/go-ipfs/blocks"
8 + u "github.com/jbenet/go-ipfs/util"
9 +)
10 +
11 +type notifications struct {
12 + wrapped *pubsub.PubSub
13 +}
14 +
15 +func newNotifications() *notifications {
16 + const bufferSize = 16
17 + return &notifications{pubsub.New(bufferSize)}
18 +}
19 +
20 +func (ps *notifications) Publish(block *blocks.Block) {
21 + topic := string(block.Key())
22 + ps.wrapped.Pub(block, topic)
23 +}
24 +
25 +// Sub returns a one-time use |blockChannel|. |blockChannel| returns nil if the
26 +// |ctx| times out or is cancelled
27 +func (ps *notifications) Subscribe(ctx context.Context, k u.Key) <-chan *blocks.Block {
28 + topic := string(k)
29 + subChan := ps.wrapped.Sub(topic)
30 + blockChannel := make(chan *blocks.Block)
31 + go func() {
32 + defer close(blockChannel)
33 + select {
34 + case val := <-subChan:
35 + block, ok := val.(*blocks.Block)
36 + if !ok {
37 + return
38 + }
39 + blockChannel <- block
40 + case <-ctx.Done():
41 + ps.wrapped.Unsub(subChan, topic)
42 + return
43 + }
44 + }()
45 + return blockChannel
46 +}
47 +
48 +func (ps *notifications) Shutdown() {
49 + ps.wrapped.Shutdown()
50 +}
bitswap/notifications_test.go new
+60
@@ -0,0 +1,60 @@
1 +package bitswap
2 +
3 +import (
4 + "bytes"
5 + "testing"
6 + "time"
7 +
8 + context "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
9 +
10 + blocks "github.com/jbenet/go-ipfs/blocks"
11 +)
12 +
13 +func TestPublishSubscribe(t *testing.T) {
14 + blockSent := getBlockOrFail(t, "Greetings from The Interval")
15 +
16 + n := newNotifications()
17 + defer n.Shutdown()
18 + ch := n.Subscribe(context.Background(), blockSent.Key())
19 +
20 + n.Publish(blockSent)
21 + blockRecvd := <-ch
22 +
23 + assertBlocksEqual(t, blockRecvd, blockSent)
24 +}
25 +
26 +func TestCarryOnWhenDeadlineExpires(t *testing.T) {
27 +
28 + impossibleDeadline := time.Nanosecond
29 + fastExpiringCtx, _ := context.WithTimeout(context.Background(), impossibleDeadline)
30 +
31 + n := newNotifications()
32 + defer n.Shutdown()
33 + blockChannel := n.Subscribe(fastExpiringCtx, getBlockOrFail(t, "A Missed Connection").Key())
34 +
35 + assertBlockChannelNil(t, blockChannel)
36 +}
37 +
38 +func assertBlockChannelNil(t *testing.T, blockChannel <-chan *blocks.Block) {
39 + blockReceived := <-blockChannel
40 + if blockReceived != nil {
41 + t.Fail()
42 + }
43 +}
44 +
45 +func assertBlocksEqual(t *testing.T, a, b *blocks.Block) {
46 + if !bytes.Equal(a.Data, b.Data) {
47 + t.Fail()
48 + }
49 + if a.Key() != b.Key() {
50 + t.Fail()
51 + }
52 +}
53 +
54 +func getBlockOrFail(t *testing.T, msg string) *blocks.Block {
55 + block, blockCreationErr := blocks.NewBlock([]byte(msg))
56 + if blockCreationErr != nil {
57 + t.Fail()
58 + }
59 + return block
60 +}