-
Notifications
You must be signed in to change notification settings - Fork 4
/
req_buf.go
59 lines (51 loc) · 930 Bytes
/
req_buf.go
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
package peer
import (
"sync"
"sync/atomic"
"github.com/ibalajiarun/go-consensus/pkg/command/commandpb"
)
const propBufCap = 1024
type reqBuf struct {
mu sync.RWMutex
full sync.Cond
b [propBufCap]reqBufElem
i int32
threshold int32
}
type reqBufElem struct {
cmd *commandpb.Command
c chan<- *commandpb.CommandResult
}
func (b *reqBuf) init(threshold int32) {
b.threshold = threshold
b.full.L = b.mu.RLocker()
}
func (b *reqBuf) add(e reqBufElem) {
b.mu.RLock()
defer b.mu.RUnlock()
for {
n := atomic.AddInt32(&b.i, 1)
i := int(n - 1)
if i < len(b.b) {
b.b[i] = e
return
}
b.full.Wait()
}
}
func (b *reqBuf) flush(f func([]reqBufElem), force bool) {
if atomic.LoadInt32(&b.i) < b.threshold && !force {
return
}
b.mu.Lock()
defer b.mu.Unlock()
i := int(b.i)
if i >= len(b.b) {
i = len(b.b)
b.full.Broadcast()
}
if i > 0 {
f(b.b[:i])
b.i = 0
}
}