added ctxgroup and router
Juan Batiz-Benet committed
Dec 13, 2014 at 04:58 UTC
62204fce6541fceb161ddbf0d8273c73647af4b5
17 files changed
+1055
Godeps/Godeps.json
+8
@@ -92,6 +92,10 @@
92
"ImportPath": "github.com/jbenet/go-base58",
93
"Rev": "568a28d73fd97651d3442392036a658b6976eed5"
94
},
95
+ {
96
+ "ImportPath": "github.com/jbenet/go-ctxgroup",
97
+ "Rev": "6b9437e8517175306e30ffec241c523752a38303"
98
+ },
99
{
100
"ImportPath": "github.com/jbenet/go-datastore",
101
"Rev": "6a1c83bda2a71a9bdc936749fdb507df958ed949"
@@ -126,6 +130,10 @@
130
"ImportPath": "github.com/jbenet/go-random",
131
"Rev": "2e83344e7dc7898f94501665af34edd4aa95a013"
132
},
133
+ {
134
+ "ImportPath": "github.com/jbenet/go-router",
135
+ "Rev": "7a4053217b7bfe3a14cc79541557d47d0c4ad85f"
136
+ },
137
{
138
"ImportPath": "github.com/kr/binarydist",
139
"Rev": "9955b0ab8708602d411341e55fffd7e0700f86bd"
Godeps/_workspace/src/github.com/jbenet/go-ctxgroup/Godeps/Godeps.json
new
+14
@@ -0,0 +1,14 @@
1
+{
2
+ "ImportPath": "github.com/jbenet/go-ctxgroup",
3
+ "GoVersion": "go1.3",
4
+ "Packages": [
5
+ "./..."
6
+ ],
7
+ "Deps": [
8
+ {
9
+ "ImportPath": "code.google.com/p/go.net/context",
10
+ "Comment": "null-144",
11
+ "Rev": "ad01a6fcc8a19d3a4478c836895ffe883bd2ceab"
12
+ }
13
+ ]
14
+}
Godeps/_workspace/src/github.com/jbenet/go-ctxgroup/Godeps/Readme
new
+5
@@ -0,0 +1,5 @@
1
+This directory tree is generated automatically by godep.
2
+
3
+Please do not edit.
4
+
5
+See https://github.com/tools/godep for more information.
Godeps/_workspace/src/github.com/jbenet/go-ctxgroup/Makefile
new
+16
@@ -0,0 +1,16 @@
1
+all:
2
+ # no-op
3
+
4
+GODEP=$(which godep)
5
+
6
+godep: ${GODEP}
7
+
8
+${GODEP}:
9
+ echo ${GODEP}
10
+ go get github.com/tools/godep
11
+
12
+# saves/vendors third-party dependencies to Godeps/_workspace
13
+# -r flag rewrites import paths to use the vendored path
14
+# ./... performs operation on all packages in tree
15
+vendor: godep
16
+ godep save -r ./...
Godeps/_workspace/src/github.com/jbenet/go-ctxgroup/README.md
new
+35
@@ -0,0 +1,35 @@
1
+# ContextGroup
2
+
3
+
4
+- Godoc: https://godoc.org/github.com/jbenet/go-ctxgroup
5
+
6
+ContextGroup is an interface for services able to be opened and closed.
7
+It has a parent Context, and Children. But ContextGroup is not a proper
8
+"tree" like the Context tree. It is more like a Context-WaitGroup hybrid.
9
+It models a main object with a few children objects -- and, unlike the
10
+context -- concerns itself with the parent-child closing semantics:
11
+
12
+- Can define an optional TeardownFunc (func() error) to be run at Closetime.
13
+- Children call Children().Add(1) to be waited upon
14
+- Children can select on <-Closing() to know when they should shut down.
15
+- Close() will wait until all children call Children().Done()
16
+- <-Closed() signals when the service is completely closed.
17
+
18
+ContextGroup can be embedded into the main object itself. In that case,
19
+the teardownFunc (if a member function) has to be set after the struct
20
+is intialized:
21
+
22
+```Go
23
+type service struct {
24
+ ContextGroup
25
+ net.Conn
26
+}
27
+func (s *service) close() error {
28
+ return s.Conn.Close()
29
+}
30
+func newService(ctx context.Context, c net.Conn) *service {
31
+ s := &service{c}
32
+ s.ContextGroup = NewContextGroup(ctx, s.close)
33
+ return s
34
+}
35
+```
Godeps/_workspace/src/github.com/jbenet/go-ctxgroup/ctxgroup.go
new
+265
@@ -0,0 +1,265 @@
1
+// package ctxgroup provides the ContextGroup, a hybrid between the
2
+// context.Context and sync.WaitGroup, which models process trees.
3
+package ctxgroup
4
+
5
+import (
6
+ "sync"
7
+
8
+ context "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
9
+)
10
+
11
+// TeardownFunc is a function used to cleanup state at the end of the
12
+// lifecycle of a process.
13
+type TeardownFunc func() error
14
+
15
+// ChildFunc is a function to register as a child. It will be automatically
16
+// tracked.
17
+type ChildFunc func(parent ContextGroup)
18
+
19
+var nilTeardownFunc = func() error { return nil }
20
+
21
+// ContextGroup is an interface for services able to be opened and closed.
22
+// It has a parent Context, and Children. But ContextGroup is not a proper
23
+// "tree" like the Context tree. It is more like a Context-WaitGroup hybrid.
24
+// It models a main object with a few children objects -- and, unlike the
25
+// context -- concerns itself with the parent-child closing semantics:
26
+//
27
+// - Can define an optional TeardownFunc (func() error) to be run at Close time.
28
+// - Children call Children().Add(1) to be waited upon
29
+// - Children can select on <-Closing() to know when they should shut down.
30
+// - Close() will wait until all children call Children().Done()
31
+// - <-Closed() signals when the service is completely closed.
32
+//
33
+// ContextGroup can be embedded into the main object itself. In that case,
34
+// the teardownFunc (if a member function) has to be set after the struct
35
+// is intialized:
36
+//
37
+// type service struct {
38
+// ContextGroup
39
+// net.Conn
40
+// }
41
+//
42
+// func (s *service) close() error {
43
+// return s.Conn.Close()
44
+// }
45
+//
46
+// func newService(ctx context.Context, c net.Conn) *service {
47
+// s := &service{c}
48
+// s.ContextGroup = NewContextGroup(ctx, s.close)
49
+// return s
50
+// }
51
+//
52
+type ContextGroup interface {
53
+
54
+ // Context is the context of this ContextGroup. It is "sort of" a parent.
55
+ Context() context.Context
56
+
57
+ // SetTeardown assigns the teardown function.
58
+ // It is called exactly _once_ when the ContextGroup is Closed.
59
+ SetTeardown(tf TeardownFunc)
60
+
61
+ // Children is a sync.Waitgroup for all children goroutines that should
62
+ // shut down completely before this service is said to be "closed".
63
+ // Follows the semantics of WaitGroup:
64
+ //
65
+ // Children().Add(1) // add one more dependent child
66
+ // Children().Done() // child signals it is done
67
+ //
68
+ // WARNING: this is deprecated and will go away soon.
69
+ Children() *sync.WaitGroup
70
+
71
+ // AddChildGroup registers a dependent ContextGroup child. The child will
72
+ // be closed when this parent is closed, and waited upon to finish. It is
73
+ // the functional equivalent of the following:
74
+ //
75
+ // parent.Children().Add(1) // add one more dependent child
76
+ // go func(parent, child ContextGroup) {
77
+ // <-parent.Closing() // wait until parent is closing
78
+ // child.Close() // signal child to close
79
+ // parent.Children().Done() // child signals it is done
80
+ // }(a, b)
81
+ //
82
+ AddChildGroup(c ContextGroup)
83
+
84
+ // AddChildFunc registers a dependent ChildFund. The child will receive
85
+ // its parent ContextGroup, and can wait on its signals. Child references
86
+ // tracked automatically. It equivalent to the following:
87
+ //
88
+ // go func(parent, child ContextGroup) {
89
+ //
90
+ // <-parent.Closing() // wait until parent is closing
91
+ // child.Close() // signal child to close
92
+ // parent.Children().Done() // child signals it is done
93
+ // }(a, b)
94
+ //
95
+ AddChildFunc(c ChildFunc)
96
+
97
+ // Close is a method to call when you wish to stop this ContextGroup
98
+ Close() error
99
+
100
+ // Closing is a signal to wait upon, like Context.Done().
101
+ // It fires when the object should be closing (but hasn't yet fully closed).
102
+ // The primary use case is for child goroutines who need to know when
103
+ // they should shut down. (equivalent to Context().Done())
104
+ Closing() <-chan struct{}
105
+
106
+ // Closed is a method to wait upon, like Context.Done().
107
+ // It fires when the entire object is fully closed.
108
+ // The primary use case is for external listeners who need to know when
109
+ // this object is completly done, and all its children closed.
110
+ Closed() <-chan struct{}
111
+}
112
+
113
+// contextGroup is a Closer with a cancellable context
114
+type contextGroup struct {
115
+ ctx context.Context
116
+ cancel context.CancelFunc
117
+
118
+ // called to run the teardown logic.
119
+ teardownFunc TeardownFunc
120
+
121
+ // closed is released once the close function is done.
122
+ closed chan struct{}
123
+
124
+ // wait group for child goroutines
125
+ children sync.WaitGroup
126
+
127
+ // sync primitive to ensure the close logic is only called once.
128
+ closeOnce sync.Once
129
+
130
+ // error to return to clients of Close().
131
+ closeErr error
132
+}
133
+
134
+// newContextGroup constructs and returns a ContextGroup. It will call
135
+// cf TeardownFunc before its Done() Wait signals fire.
136
+func newContextGroup(ctx context.Context, cf TeardownFunc) ContextGroup {
137
+ ctx, cancel := context.WithCancel(ctx)
138
+ c := &contextGroup{
139
+ ctx: ctx,
140
+ cancel: cancel,
141
+ closed: make(chan struct{}),
142
+ }
143
+ c.SetTeardown(cf)
144
+
145
+ c.Children().Add(1) // initialize with 1. calling Close will decrement it.
146
+ go c.closeOnContextDone()
147
+ return c
148
+}
149
+
150
+// SetTeardown assigns the teardown function.
151
+func (c *contextGroup) SetTeardown(cf TeardownFunc) {
152
+ if cf == nil {
153
+ cf = nilTeardownFunc
154
+ }
155
+ c.teardownFunc = cf
156
+}
157
+
158
+func (c *contextGroup) Context() context.Context {
159
+ return c.ctx
160
+}
161
+
162
+func (c *contextGroup) Children() *sync.WaitGroup {
163
+ return &c.children
164
+}
165
+
166
+func (c *contextGroup) AddChildGroup(child ContextGroup) {
167
+ c.children.Add(1)
168
+ go func(parent, child ContextGroup) {
169
+ <-parent.Closing() // wait until parent is closing
170
+ child.Close() // signal child to close
171
+ parent.Children().Done() // child signals it is done
172
+ }(c, child)
173
+}
174
+
175
+func (c *contextGroup) AddChildFunc(child ChildFunc) {
176
+ c.children.Add(1)
177
+ go func(parent ContextGroup, child ChildFunc) {
178
+ child(parent)
179
+ parent.Children().Done() // child signals it is done
180
+ }(c, child)
181
+}
182
+
183
+// Close is the external close function. it's a wrapper around internalClose
184
+// that waits on Closed()
185
+func (c *contextGroup) Close() error {
186
+ c.internalClose()
187
+ <-c.Closed() // wait until we're totally done.
188
+ return c.closeErr
189
+}
190
+
191
+func (c *contextGroup) Closing() <-chan struct{} {
192
+ return c.Context().Done()
193
+}
194
+
195
+func (c *contextGroup) Closed() <-chan struct{} {
196
+ return c.closed
197
+}
198
+
199
+func (c *contextGroup) internalClose() {
200
+ go c.closeOnce.Do(c.closeLogic)
201
+}
202
+
203
+// the _actual_ close process.
204
+func (c *contextGroup) closeLogic() {
205
+ // this function should only be called once (hence the sync.Once).
206
+ // and it will panic at the bottom (on close(c.closed)) otherwise.
207
+
208
+ c.cancel() // signal that we're shutting down (Closing)
209
+ c.closeErr = c.teardownFunc() // actually run the close logic
210
+ c.children.Wait() // wait till all children are done.
211
+ close(c.closed) // signal that we're shut down (Closed)
212
+}
213
+
214
+// if parent context is shut down before we call Close explicitly,
215
+// we need to go through the Close motions anyway. Hence all the sync
216
+// stuff all over the place...
217
+func (c *contextGroup) closeOnContextDone() {
218
+ <-c.Context().Done() // wait until parent (context) is done.
219
+ c.internalClose()
220
+ c.Children().Done()
221
+}
222
+
223
+// WithTeardown constructs and returns a ContextGroup with
224
+// cf TeardownFunc (and context.Background)
225
+func WithTeardown(cf TeardownFunc) ContextGroup {
226
+ if cf == nil {
227
+ panic("nil TeardownFunc")
228
+ }
229
+ return newContextGroup(context.Background(), cf)
230
+}
231
+
232
+// WithContext constructs and returns a ContextGroup with given context
233
+func WithContext(ctx context.Context) ContextGroup {
234
+ if ctx == nil {
235
+ panic("nil Context")
236
+ }
237
+ return newContextGroup(ctx, nil)
238
+}
239
+
240
+// WithContextAndTeardown constructs and returns a ContextGroup with
241
+// cf TeardownFunc (and context.Background)
242
+func WithContextAndTeardown(ctx context.Context, cf TeardownFunc) ContextGroup {
243
+ if ctx == nil {
244
+ panic("nil Context")
245
+ }
246
+ if cf == nil {
247
+ panic("nil TeardownFunc")
248
+ }
249
+ return newContextGroup(ctx, cf)
250
+}
251
+
252
+// WithParent constructs and returns a ContextGroup with given parent
253
+func WithParent(p ContextGroup) ContextGroup {
254
+ if p == nil {
255
+ panic("nil ContextGroup")
256
+ }
257
+ c := newContextGroup(p.Context(), nil)
258
+ p.AddChildGroup(c)
259
+ return c
260
+}
261
+
262
+// WithBackground returns a ContextGroup with context.Background()
263
+func WithBackground() ContextGroup {
264
+ return newContextGroup(context.Background(), nil)
265
+}
Godeps/_workspace/src/github.com/jbenet/go-ctxgroup/ctxgroup_test.go
new
+187
@@ -0,0 +1,187 @@
1
+package ctxgroup
2
+
3
+import (
4
+ "testing"
5
+ "time"
6
+
7
+ context "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
8
+)
9
+
10
+type tree struct {
11
+ ContextGroup
12
+ c []tree
13
+}
14
+
15
+func setupCGHierarchy(ctx context.Context) tree {
16
+ t := func(n ContextGroup, ts ...tree) tree {
17
+ return tree{n, ts}
18
+ }
19
+
20
+ if ctx == nil {
21
+ ctx = context.Background()
22
+ }
23
+ a := WithContext(ctx)
24
+ b1 := WithParent(a)
25
+ b2 := WithParent(a)
26
+ c1 := WithParent(b1)
27
+ c2 := WithParent(b1)
28
+ c3 := WithParent(b2)
29
+ c4 := WithParent(b2)
30
+
31
+ return t(a, t(b1, t(c1), t(c2)), t(b2, t(c3), t(c4)))
32
+}
33
+
34
+func TestClosingClosed(t *testing.T) {
35
+
36
+ a := WithBackground()
37
+ Q := make(chan string)
38
+
39
+ go func() {
40
+ <-a.Closing()
41
+ Q <- "closing"
42
+ }()
43
+
44
+ go func() {
45
+ <-a.Closed()
46
+ Q <- "closed"
47
+ }()
48
+
49
+ go func() {
50
+ a.Close()
51
+ Q <- "closed"
52
+ }()
53
+
54
+ if q := <-Q; q != "closing" {
55
+ t.Error("order incorrect. closing not first")
56
+ }
57
+ if q := <-Q; q != "closed" {
58
+ t.Error("order incorrect. closing not first")
59
+ }
60
+ if q := <-Q; q != "closed" {
61
+ t.Error("order incorrect. closing not first")
62
+ }
63
+}
64
+
65
+func TestChildFunc(t *testing.T) {
66
+ a := WithBackground()
67
+
68
+ wait1 := make(chan struct{})
69
+ wait2 := make(chan struct{})
70
+ wait3 := make(chan struct{})
71
+ wait4 := make(chan struct{})
72
+ go func() {
73
+ a.Close()
74
+ wait4 <- struct{}{}
75
+ }()
76
+
77
+ a.AddChildFunc(func(parent ContextGroup) {
78
+ wait1 <- struct{}{}
79
+ <-wait2
80
+ wait3 <- struct{}{}
81
+ })
82
+
83
+ <-wait1
84
+ select {
85
+ case <-wait3:
86
+ t.Error("should not be closed yet")
87
+ case <-wait4:
88
+ t.Error("should not be closed yet")
89
+ case <-a.Closed():
90
+ t.Error("should not be closed yet")
91
+ default:
92
+ }
93
+
94
+ wait2 <- struct{}{}
95
+
96
+ select {
97
+ case <-wait3:
98
+ case <-time.After(time.Second):
99
+ t.Error("should be closed now")
100
+ }
101
+
102
+ select {
103
+ case <-wait4:
104
+ case <-time.After(time.Second):
105
+ t.Error("should be closed now")
106
+ }
107
+}
108
+
109
+func TestTeardownCalledOnce(t *testing.T) {
110
+ a := setupCGHierarchy(nil)
111
+
112
+ onlyOnce := func() func() error {
113
+ count := 0
114
+ return func() error {
115
+ count++
116
+ if count > 1 {
117
+ t.Error("called", count, "times")
118
+ }
119
+ return nil
120
+ }
121
+ }
122
+
123
+ a.SetTeardown(onlyOnce())
124
+ a.c[0].SetTeardown(onlyOnce())
125
+ a.c[0].c[0].SetTeardown(onlyOnce())
126
+ a.c[0].c[1].SetTeardown(onlyOnce())
127
+ a.c[1].SetTeardown(onlyOnce())
128
+ a.c[1].c[0].SetTeardown(onlyOnce())
129
+ a.c[1].c[1].SetTeardown(onlyOnce())
130
+
131
+ a.c[0].c[0].Close()
132
+ a.c[0].c[0].Close()
133
+ a.c[0].c[0].Close()
134
+ a.c[0].c[0].Close()
135
+ a.c[0].Close()
136
+ a.c[0].Close()
137
+ a.c[0].Close()
138
+ a.c[0].Close()
139
+ a.Close()
140
+ a.Close()
141
+ a.Close()
142
+ a.Close()
143
+ a.c[1].Close()
144
+ a.c[1].Close()
145
+ a.c[1].Close()
146
+ a.c[1].Close()
147
+}
148
+
149
+func TestOnClosed(t *testing.T) {
150
+
151
+ ctx, cancel := context.WithCancel(context.Background())
152
+ a := setupCGHierarchy(ctx)
153
+ Q := make(chan string, 10)
154
+
155
+ onClosed := func(s string, c ContextGroup) {
156
+ <-c.Closed()
157
+ Q <- s
158
+ }
159
+
160
+ go onClosed("0", a.c[0])
161
+ go onClosed("10", a.c[1].c[0])
162
+ go onClosed("", a)
163
+ go onClosed("00", a.c[0].c[0])
164
+ go onClosed("1", a.c[1])
165
+ go onClosed("01", a.c[0].c[1])
166
+ go onClosed("11", a.c[1].c[1])
167
+
168
+ test := func(ss ...string) {
169
+ s1 := <-Q
170
+ for _, s2 := range ss {
171
+ if s1 == s2 {
172
+ return
173
+ }
174
+ }
175
+ t.Error("context not in group", s1, ss)
176
+ }
177
+
178
+ cancel()
179
+
180
+ test("00", "01", "10", "11")
181
+ test("00", "01", "10", "11")
182
+ test("00", "01", "10", "11")
183
+ test("00", "01", "10", "11")
184
+ test("0", "1")
185
+ test("0", "1")
186
+ test("")
187
+}
Godeps/_workspace/src/github.com/jbenet/go-router/LICENSE
new
+21
@@ -0,0 +1,21 @@
1
+The MIT License (MIT)
2
+
3
+Copyright (c) 2014 Juan Batiz-Benet
4
+
5
+Permission is hereby granted, free of charge, to any person obtaining a copy
6
+of this software and associated documentation files (the "Software"), to deal
7
+in the Software without restriction, including without limitation the rights
8
+to use, copy, modify, merge, publish, distribute, sublicense, and/or sell
9
+copies of the Software, and to permit persons to whom the Software is
10
+furnished to do so, subject to the following conditions:
11
+
12
+The above copyright notice and this permission notice shall be included in
13
+all copies or substantial portions of the Software.
14
+
15
+THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR
16
+IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY,
17
+FITNESS FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE
18
+AUTHORS OR COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER
19
+LIABILITY, WHETHER IN AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM,
20
+OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN
21
+THE SOFTWARE.
Godeps/_workspace/src/github.com/jbenet/go-router/README.md
new
+8
@@ -0,0 +1,8 @@
1
+# go-router - networking inspired muxing
2
+
3
+This is a networking-inspired model. It's useful for muxing and lots of other
4
+things. It enables you to construct your abstractions as you would compose
5
+computer networks (endpoints, switches, routing tables). These could represent
6
+processing workers, entire subsystems, or even real computers ;).
7
+
8
+See https://godoc.org/github.com/jbenet/go-router
Godeps/_workspace/src/github.com/jbenet/go-router/distance.go
new
+33
@@ -0,0 +1,33 @@
1
+package router
2
+
3
+// HammingDistance is a DistanceFunc that interprets Addresses as strings
4
+// and uses their Hamming distance.
5
+// Return -1 if the Addresses are not strings, or strings length don't match.
6
+func HammingDistance(a1, a2 Address) int {
7
+ s1, ok := a1.(string)
8
+ if !ok {
9
+ return -1
10
+ }
11
+
12
+ s2, ok := a2.(string)
13
+ if !ok {
14
+ return -1
15
+ }
16
+
17
+ // runes not code points
18
+ r1 := []rune(s1)
19
+ r2 := []rune(s2)
20
+
21
+ // hamming distance requires equal length strings
22
+ if len(r1) != len(r2) {
23
+ return -1
24
+ }
25
+
26
+ d := 0
27
+ for i := range r1 {
28
+ if r1[i] != r2[i] {
29
+ d++
30
+ }
31
+ }
32
+ return d
33
+}
Godeps/_workspace/src/github.com/jbenet/go-router/impl.go
new
+81
@@ -0,0 +1,81 @@
1
+package router
2
+
3
+import (
4
+ "errors"
5
+)
6
+
7
+type packet struct {
8
+ a Address
9
+ p interface{}
10
+}
11
+
12
+// ErrNoRoute signals when there is no Route to a destination
13
+var ErrNoRoute = errors.New("routing error: no route")
14
+
15
+// NewPacket constructs a trivial packet linking a destination Address to
16
+// an interface{} payload.
17
+func NewPacket(destination Address, payload interface{}) Packet {
18
+ return &packet{destination, payload}
19
+}
20
+
21
+func (p *packet) Destination() Address {
22
+ return p.a
23
+}
24
+
25
+func (p *packet) Payload() interface{} {
26
+ return p.p
27
+}
28
+
29
+// QueueNode is a trivial node, which accepts packets into a queue
30
+type QueueNode struct {
31
+ a Address
32
+ q chan Packet
33
+}
34
+
35
+// NewQueueNode constructs a node with an internal chan Packet queue
36
+func NewQueueNode(addr Address, q chan Packet) *QueueNode {
37
+ return &QueueNode{addr, q}
38
+}
39
+
40
+// Queue returns the chan Packet queue
41
+func (n *QueueNode) Queue() <-chan Packet {
42
+ return n.q
43
+}
44
+
45
+// Address returns the QueueNode's Address
46
+func (n *QueueNode) Address() Address {
47
+ return n.a
48
+}
49
+
50
+// HandlePacket consumes the incomng packet and adds it to the queue.
51
+func (n *QueueNode) HandlePacket(p Packet, s Node) error {
52
+ n.q <- p
53
+ return nil
54
+}
55
+
56
+type switchh struct {
57
+ addr Address
58
+ router Router
59
+ nodes []Node
60
+}
61
+
62
+// NewSwitch constructs a switch with given Router and list of adjacent Nodes.
63
+func NewSwitch(a Address, r Router, adj []Node) Switch {
64
+ return &switchh{a, r, adj}
65
+}
66
+
67
+func (s *switchh) Address() Address {
68
+ return s.addr
69
+}
70
+
71
+func (s *switchh) Router() Router {
72
+ return s.router
73
+}
74
+
75
+func (s *switchh) HandlePacket(p Packet, n Node) error {
76
+ next := s.router.Route(p)
77
+ if next != nil {
78
+ return next.HandlePacket(p, s)
79
+ }
80
+ return ErrNoRoute
81
+}
Godeps/_workspace/src/github.com/jbenet/go-router/interface.go
new
+81
@@ -0,0 +1,81 @@
1
+// Package router is a networking-inspired routing model. It's useful for
2
+// muxing and lots of other things. It enables you to construct your
3
+// abstractions as you would compose computer networks (endpoints, switches,
4
+// routing tables). These could represent processing workers, entire
5
+// subsystems, or even real computers ;).
6
+package router
7
+
8
+// Address is our way of knowing where we're headed. Traditionally, addresses
9
+// are things like IP Addresses, or email addresses. But filepaths, or URLs
10
+// can be seen as addresses too. We really leave it up to you. Our routing
11
+// is general, and forces you to pick an addressing scheme, and some logic to
12
+// discriminate addresses that you'll plug into Routers (Address DistanceFunc).
13
+type Address interface{}
14
+
15
+// Packet is the unit of moving things in our network. Anything can be routed
16
+// in our network, as long as it has a Destination.
17
+type Packet interface {
18
+
19
+ // Destination is the Address of the endpoint this Packet is headed to.
20
+ // They could use the same Addressing throughout (recommended) like the
21
+ // internet, or face the pain of translating addresses at inter-network
22
+ // gateways.
23
+ Destination() Address
24
+
25
+ // Payload is here for completeness. This can be the Packet itself, but this
26
+ // function encourages the client to think through their packet design.
27
+ Payload() interface{}
28
+}
29
+
30
+// Node is an object which has interfaces to connect to networks. This
31
+// is an "endpoint" object.
32
+type Node interface {
33
+
34
+ // Address returns the node's address.
35
+ Address() Address
36
+
37
+ // HandlePacket receives a packet sent by another node.
38
+ HandlePacket(Packet, Node) error
39
+}
40
+
41
+// Switch is the basic forwarding device. It listens to all its interfaces (it
42
+// is a Node in our network) and Forwards all Packets received by any interface
43
+// out another interface, according to its ForwardingTable.
44
+type Switch interface {
45
+ Node
46
+
47
+ // Router returns the Router used to decide where to forward packets.
48
+ Router() Router
49
+}
50
+
51
+// Router is an object that decides how a Packet should be Routed. This should
52
+// be as close to a static table lookup as possible, meaning it would be best
53
+// to prepare a forwarding table in parallel, instead of blocking.
54
+//
55
+// Our Router captures the entire Control Plane, meaning that we can implement:
56
+// - Static Routing - forwarding table only
57
+// - Dynamic Routing - routing table computed with an algorithm or protocol
58
+// - "SDN" Routing - Control Plane separated from Data Plane
59
+// And even:
60
+// - URL Routers (like gorilla.Muxer)
61
+// - Protocol Muxers
62
+// entirely within different Router implementations.
63
+//
64
+// Note that this is a break from traditional networking systems. Instead of
65
+// having the abstractions of FIB, RIB, Routing/Forwarding Tables, Routers,
66
+// and Switches, we only have the last two:
67
+// - Router -- the things that "route" (decide where things go)
68
+// - Switch -- connecting Nodes, "switch" Packets according to a Router.
69
+type Router interface {
70
+
71
+ // Route decides how to route a Packet.
72
+ // It returns the next hop Node chosen to send the Packet to.
73
+ // Route may return nil, if no route is suitable at all (equivalent of drop).
74
+ Route(Packet) Node
75
+}
76
+
77
+// DistanceFunc returns a measure of distance between two Addresses.
78
+// Examples:
79
+// - masked longest prefix match (IP, CIDR)
80
+// - XOR distance (Kademlia)
81
+type DistanceFunc func(a, b Address) int
Godeps/_workspace/src/github.com/jbenet/go-router/ip/ip.go
new
+6
@@ -0,0 +1,6 @@
1
+// Package ip provides a Table for ip matching
2
+package ip
3
+
4
+// import (
5
+// ipaddr "github.com/mikioh/ipaddr"
6
+// )
Godeps/_workspace/src/github.com/jbenet/go-router/ip/ip_test.go
new
+1
@@ -0,0 +1 @@
1
+package ip
Godeps/_workspace/src/github.com/jbenet/go-router/pkt_test.go
new
+77
@@ -0,0 +1,77 @@
1
+package router
2
+
3
+import (
4
+ "testing"
5
+)
6
+
7
+type mockNode struct {
8
+ addr string
9
+ pkts []Packet
10
+}
11
+
12
+func (m *mockNode) Address() Address {
13
+ return m.addr
14
+}
15
+
16
+func (m *mockNode) HandlePacket(p Packet, n Node) error {
17
+ m.pkts = append(m.pkts, p)
18
+ return nil
19
+}
20
+
21
+type mockPacket struct {
22
+ a string
23
+}
24
+
25
+func (m *mockPacket) Destination() Address {
26
+ return m.a
27
+}
28
+
29
+func (m *mockPacket) Payload() interface{} {
30
+ return m
31
+}
32
+
33
+func TestAddrs(t *testing.T) {
34
+
35
+ ta := func(a, b Address, expect int) {
36
+ actual := HammingDistance(a, b)
37
+ if actual != expect {
38
+ t.Error("address distance error:", a, b, expect, actual)
39
+ }
40
+ }
41
+
42
+ a := "abc"
43
+ b := "abc"
44
+ c := "abd"
45
+ d := "add"
46
+ e := "ddd"
47
+
48
+ ta(a, a, 0)
49
+ ta(a, b, 0)
50
+ ta(a, c, 1)
51
+ ta(a, d, 2)
52
+ ta(a, e, 3)
53
+}
54
+
55
+func TestNodes(t *testing.T) {
56
+
57
+ a := &mockPacket{"abc"}
58
+ b := &mockPacket{"abc"}
59
+ c := &mockPacket{"abd"}
60
+ d := &mockPacket{"add"}
61
+ e := &mockPacket{"ddd"}
62
+
63
+ n := &mockNode{addr: "abc"}
64
+ n2 := &mockNode{addr: "ddd"}
65
+ n.HandlePacket(a, n2)
66
+ n.HandlePacket(b, n2)
67
+ n.HandlePacket(c, n2)
68
+ n.HandlePacket(d, n2)
69
+ n.HandlePacket(e, n2)
70
+
71
+ pkts := []Packet{a, b, c, d, e}
72
+ for i, p := range pkts {
73
+ if n.pkts[i] != p {
74
+ t.Error("pkts not handled in order.")
75
+ }
76
+ }
77
+}
Godeps/_workspace/src/github.com/jbenet/go-router/table.go
new
+137
@@ -0,0 +1,137 @@
1
+package router
2
+
3
+import (
4
+ "sync"
5
+)
6
+
7
+// TableEntry is an (Address, Node) pair for a routing Table
8
+type TableEntry interface {
9
+ Address() Address
10
+ NextHop() Node
11
+}
12
+
13
+func NewTableEntry(addr Address, nexthop Node) TableEntry {
14
+ return &tableEntry{addr, nexthop}
15
+}
16
+
17
+type tableEntry struct {
18
+ addr Address
19
+ next Node
20
+}
21
+
22
+func (te *tableEntry) Address() Address {
23
+ return te.addr
24
+}
25
+
26
+func (te *tableEntry) NextHop() Node {
27
+ return te.next
28
+}
29
+
30
+// Table is a Router (Routing Table, really) based on a distance criterion.
31
+//
32
+// For example:
33
+//
34
+// n1 := NewQueueNode("aaa", make(chan Packet, 10))
35
+// n2 := NewQueueNode("aba", make(chan Packet, 10))
36
+// n3 := NewQueueNode("abc", make(chan Packet, 10))
37
+//
38
+// var t router.Table
39
+// t.Distance = router.HammingDistance
40
+// t.AddNodes(n1, n2)
41
+//
42
+// p1 := NewPacket("aaa", "hello1")
43
+// p2 := NewPacket("aba", "hello2")
44
+// p3 := NewPacket("abc", "hello3")
45
+//
46
+// t.Route(p1) // n1
47
+// t.Route(p2) // n2
48
+// t.Route(p3) // n2, because we don't have n3 and n2 is closet
49
+//
50
+// t.AddNode(n3)
51
+// t.Route(p3) // n3
52
+type Table interface {
53
+ Router
54
+
55
+ // Entries are the entries in this routing table
56
+ Entries() []TableEntry
57
+
58
+ // Distance returns a measure of distance between two Addresses
59
+ Distance() DistanceFunc
60
+}
61
+
62
+type SimpleTable struct {
63
+ entries []TableEntry
64
+ distance DistanceFunc
65
+ sync.RWMutex
66
+}
67
+
68
+// Entries are the entries in this routing table
69
+func (t *SimpleTable) Entries() []TableEntry {
70
+ return t.entries
71
+}
72
+
73
+// Distance returns a measure of distance between two Addresses.
74
+func (t *SimpleTable) Distance() DistanceFunc {
75
+ return t.distance
76
+}
77
+
78
+// AddEntry adds an (Address, NextHop) entry to the Table
79
+func (t *SimpleTable) AddEntry(addr Address, nextHop Node) {
80
+ t.Lock()
81
+ defer t.Unlock()
82
+ t.entries = append(t.entries, NewTableEntry(addr, nextHop))
83
+}
84
+
85
+// AddNode calls AddTableEntry for the given Node
86
+func (t *SimpleTable) AddNode(n Node) {
87
+ t.AddEntry(n.Address(), n)
88
+}
89
+
90
+// AddNodes calls AddTableEntry for the given Node
91
+func (t *SimpleTable) AddNodes(ns ...Node) {
92
+ t.Lock()
93
+ defer t.Unlock()
94
+ for _, n := range ns {
95
+ t.entries = append(t.entries, NewTableEntry(n.Address(), n))
96
+ }
97
+}
98
+
99
+// Route decides how to route a Packet out of a list of Nodes.
100
+// It returns the Node chosen to send the Packet to.
101
+// Route may return nil, if no route is suitable at all (equivalent of drop).
102
+func (t *SimpleTable) Route(p Packet) Node {
103
+ if t.entries == nil {
104
+ return nil
105
+ }
106
+
107
+ t.RLock()
108
+ defer t.RUnlock()
109
+
110
+ dist := t.Distance()
111
+ if dist == nil {
112
+ dist = equalDistance
113
+ }
114
+
115
+ var best Node
116
+ var bestDist int
117
+ var addr = p.Destination()
118
+
119
+ for _, e := range t.entries {
120
+ d := dist(e.Address(), addr)
121
+ if d < 0 {
122
+ continue
123
+ }
124
+ if best == nil || d < bestDist {
125
+ bestDist = d
126
+ best = e.NextHop()
127
+ }
128
+ }
129
+ return best
130
+}
131
+
132
+func equalDistance(a, b Address) int {
133
+ if a == b {
134
+ return 0
135
+ }
136
+ return -1
137
+}
Godeps/_workspace/src/github.com/jbenet/go-router/table_test.go
new
+80
@@ -0,0 +1,80 @@
1
+package router
2
+
3
+import (
4
+ "testing"
5
+)
6
+
7
+func TestTable(t *testing.T) {
8
+
9
+ na := &mockNode{addr: "abc"}
10
+ nc := &mockNode{addr: "abd"}
11
+ nd := &mockNode{addr: "add"}
12
+ ne := &mockNode{addr: "ddd"}
13
+
14
+ pa := &mockPacket{"abc"}
15
+ pb := &mockPacket{"abc"}
16
+ pc := &mockPacket{"abd"}
17
+ pd := &mockPacket{"add"}
18
+ pe := &mockPacket{"ddd"}
19
+
20
+ table := &SimpleTable{
21
+ entries: []TableEntry{
22
+ &tableEntry{na.Address(), na},
23
+ &tableEntry{nc.Address(), nc},
24
+ &tableEntry{nd.Address(), nd},
25
+ &tableEntry{ne.Address(), ne},
26
+ },
27
+ }
28
+
29
+ s := NewSwitch("sss", table, []Node{na, nc, nd, ne})
30
+ s.HandlePacket(pa, na)
31
+ s.HandlePacket(pb, na)
32
+ s.HandlePacket(pc, na)
33
+ s.HandlePacket(pd, na)
34
+ s.HandlePacket(pe, na)
35
+
36
+ tt := func(n *mockNode, pkts []Packet) {
37
+ for i, p := range pkts {
38
+ if len(n.pkts) <= i {
39
+ t.Error("pkts not handled in order.", n, pkts)
40
+ return
41
+ }
42
+ if n.pkts[i] != p {
43
+ t.Error("pkts not handled in order.", n, pkts)
44
+ }
45
+ }
46
+ }
47
+
48
+ tt(na, []Packet{pa, pb})
49
+ tt(nc, []Packet{pc})
50
+ tt(nd, []Packet{pd})
51
+ tt(ne, []Packet{pe})
52
+}
53
+
54
+func TestTable2(t *testing.T) {
55
+
56
+ n1 := NewQueueNode("aaa", make(chan Packet, 1))
57
+ n2 := NewQueueNode("aba", make(chan Packet, 1))
58
+ n3 := NewQueueNode("abc", make(chan Packet, 1))
59
+
60
+ var tb SimpleTable
61
+ tb.distance = HammingDistance
62
+ tb.AddNodes(n1, n2)
63
+
64
+ p1 := NewPacket("aaa", "hello1")
65
+ p2 := NewPacket("aba", "hello2")
66
+ p3 := NewPacket("abc", "hello3")
67
+
68
+ testRoute := func(p Packet, expect Node) {
69
+ if tb.Route(p) != expect {
70
+ t.Error(p, "route should be", expect)
71
+ }
72
+ }
73
+
74
+ testRoute(p1, n1)
75
+ testRoute(p2, n2)
76
+ testRoute(p3, n2) // n2 because we don't have n3 and n2 is closet
77
+
78
+ tb.AddNode(n3)
79
+ testRoute(p3, n3)
80
+}