-
Notifications
You must be signed in to change notification settings - Fork 3
/
put_buf.go
63 lines (55 loc) · 1.13 KB
/
put_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
60
61
62
63
// Copyright (C) 2018 Kun Zhong All rights reserved.
// Use of this source code is governed by a Licensed under the Apache License, Version 2.0 (the "License");
package grpcx
import (
"bytes"
"reflect"
"sync"
)
// putBuffer is an unbounded channel of outPack structs.
type outPack struct {
mtype reflect.Type
body *bytes.Buffer
}
type putBuffer struct {
c chan *outPack
mu sync.Mutex
backlog []*outPack
}
func newPutBuffer() *putBuffer {
b := &putBuffer{
c: make(chan *outPack, 1),
}
return b
}
func (b *putBuffer) put(r *outPack) {
b.mu.Lock()
defer b.mu.Unlock()
if len(b.backlog) == 0 {
select {
case b.c <- r:
return
default:
}
}
b.backlog = append(b.backlog, r)
}
func (b *putBuffer) load() {
b.mu.Lock()
defer b.mu.Unlock()
if len(b.backlog) > 0 {
select {
case b.c <- b.backlog[0]:
b.backlog[0] = nil
b.backlog = b.backlog[1:]
default:
}
}
}
// get returns the channel that put a outPack in the buffer.
//
// Upon receipt of a outPack, the caller should call load to send another
// outPack onto the channel if there is any.
func (b *putBuffer) get() <-chan *outPack {
return b.c
}