updated msgio (bufixes)
Juan Batiz-Benet committed
Dec 13, 2014 at 07:12 UTC
694453102691971435bbacf59acc26bfbe8576cf
5 files changed
+89
-19
Godeps/Godeps.json
+1
-1
@@ -110,7 +110,7 @@
110
},
111
{
112
"ImportPath": "github.com/jbenet/go-msgio",
113
- "Rev": "7bdc5b738564871e1c0d5ca9449900d0d6773713"
113
+ "Rev": "753e598a1d24b311ee05c4ce001cff74e2a8e745"
114
},
115
{
116
"ImportPath": "github.com/jbenet/go-multiaddr",
Godeps/_workspace/src/github.com/jbenet/go-msgio/chan.go
+1
-1
@@ -3,7 +3,7 @@ package msgio
3
import (
4
"io"
5
6
- mpool "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-msgio/mpool"
6
+ mpool "github.com/jbenet/go-msgio/mpool"
7
)
8
9
// Chan is a msgio duplex channel. It is used to have a channel interface
Godeps/_workspace/src/github.com/jbenet/go-msgio/mpool/pool.go
+31
-15
@@ -72,8 +72,9 @@ func (p *Pool) getPool(idx uint32) *sync.Pool {
72
// If Get would otherwise return nil and p.New is non-nil, Get returns the
73
// result of calling p.New.
74
func (p *Pool) Get(length uint32) interface{} {
75
- idx := largerPowerOfTwo(length)
75
+ idx := nextPowerOfTwo(length)
76
sp := p.getPool(idx)
77
+ // fmt.Printf("Get(%d) idx(%d)\n", length, idx)
78
val := sp.Get()
79
if val == nil && p.New != nil {
80
val = p.New(0x1 << idx)
@@ -83,27 +84,42 @@ func (p *Pool) Get(length uint32) interface{} {
84
85
// Put adds x to the pool.
86
func (p *Pool) Put(length uint32, val interface{}) {
86
- idx := smallerPowerOfTwo(length)
87
+ idx := prevPowerOfTwo(length)
88
+ // fmt.Printf("Put(%d, -) idx(%d)\n", length, idx)
89
sp := p.getPool(idx)
90
sp.Put(val)
91
}
92
91
-func largerPowerOfTwo(num uint32) uint32 {
92
- for p := uint32(0); p < 32; p++ {
93
- if (0x1 << p) >= num {
94
- return p
95
- }
93
+func nextPowerOfTwo(v uint32) uint32 {
94
+ // fmt.Printf("nextPowerOfTwo(%d) ", v)
95
+ v--
96
+ v |= v >> 1
97
+ v |= v >> 2
98
+ v |= v >> 4
99
+ v |= v >> 8
100
+ v |= v >> 16
101
+ v++
102
+
103
+ // fmt.Printf("-> %d", v)
104
+
105
+ i := uint32(0)
106
+ for i = 0; v > 1; i++ {
107
+ v = v >> 1
108
}
109
98
- panic("unreachable")
110
+ // fmt.Printf("-> %d\n", i)
111
+ return i
112
}
113
101
-func smallerPowerOfTwo(num uint32) uint32 {
102
- for p := uint32(1); p < 32; p++ {
103
- if (0x1 << p) > num {
104
- return p - 1
105
- }
114
+func prevPowerOfTwo(num uint32) uint32 {
115
+ next := nextPowerOfTwo(num)
116
+ // fmt.Printf("prevPowerOfTwo(%d) next: %d", num, next)
117
+ switch {
118
+ case num == (1 << next): // num is a power of 2
119
+ case next == 0:
120
+ default:
121
+ next = next - 1 // smaller
122
}
107
-
108
- panic("unreachable")
123
+ // fmt.Printf(" = %d\n", next)
124
+ return next
125
}
Godeps/_workspace/src/github.com/jbenet/go-msgio/mpool/pool_test.go
+55
-1
@@ -9,6 +9,7 @@ package mpool
9
10
import (
11
"fmt"
12
+ "math/rand"
13
"runtime"
14
"runtime/debug"
15
"sync/atomic"
@@ -62,7 +63,7 @@ func TestPoolNew(t *testing.T) {
63
s := [32]int{}
64
p := Pool{
65
New: func(length int) interface{} {
65
- idx := largerPowerOfTwo(uint32(length))
66
+ idx := nextPowerOfTwo(uint32(length))
67
s[idx]++
68
return s[idx]
69
},
@@ -148,10 +149,63 @@ func TestPoolStress(t *testing.T) {
149
}()
150
}
151
for i := 0; i < P; i++ {
152
+ // fmt.Printf("%d/%d\n", i, P)
153
<-done
154
}
155
}
156
157
+func TestPoolStressByteSlicePool(t *testing.T) {
158
+ const P = 10
159
+ chs := 10
160
+ maxSize := uint32(1 << 16)
161
+ N := int(1e4)
162
+ if testing.Short() {
163
+ N /= 100
164
+ }
165
+ p := ByteSlicePool
166
+ done := make(chan bool)
167
+ errs := make(chan error)
168
+ for i := 0; i < P; i++ {
169
+ go func() {
170
+ ch := make(chan []byte, chs+1)
171
+
172
+ for i := 0; i < chs; i++ {
173
+ j := rand.Uint32() % maxSize
174
+ ch <- p.Get(j).([]byte)
175
+ }
176
+
177
+ for j := 0; j < N; j++ {
178
+ r := uint32(0)
179
+ for i := 0; i < chs; i++ {
180
+ v := <-ch
181
+ p.Put(uint32(cap(v)), v)
182
+ r = rand.Uint32() % maxSize
183
+ v = p.Get(r).([]byte)
184
+ if uint32(len(v)) < r {
185
+ errs <- fmt.Errorf("expect len(v) >= %d, got %d", j, len(v))
186
+ }
187
+ ch <- v
188
+ }
189
+
190
+ if r%1000 == 0 {
191
+ runtime.GC()
192
+ }
193
+ }
194
+ done <- true
195
+ }()
196
+ }
197
+
198
+ for i := 0; i < P; {
199
+ select {
200
+ case <-done:
201
+ i++
202
+ // fmt.Printf("%d/%d\n", i, P)
203
+ case err := <-errs:
204
+ t.Error(err)
205
+ }
206
+ }
207
+}
208
+
209
func BenchmarkPool(b *testing.B) {
210
var p Pool
211
b.RunParallel(func(pb *testing.PB) {
Godeps/_workspace/src/github.com/jbenet/go-msgio/msgio.go
+1
-1
@@ -5,7 +5,7 @@ import (
5
"io"
6
"sync"
7
8
- mpool "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-msgio/mpool"
8
+ mpool "github.com/jbenet/go-msgio/mpool"
9
)
10
11
// NBO is NetworkByteOrder