-
Notifications
You must be signed in to change notification settings - Fork 1
/
rcache.go
155 lines (137 loc) · 3.07 KB
/
rcache.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
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
package collect
import (
"context"
"hash/fnv"
lru "github.com/hashicorp/golang-lru"
)
// requestCache .
// requestCache is used to store the request control message,
// which is for response routing.
// TODO: rename it to requestCache.
type requestWorkerPool struct {
cache *lru.Cache
}
// reqstItem .
type reqItem struct {
finalHandler FinalRespHandler
topic string
req *Request
}
type reqWorker struct {
item *reqItem
cc func()
onClose func()
}
func newReqWorker(ctx context.Context, item *reqItem, onClose func()) *reqWorker {
cctx, cc := context.WithCancel(ctx)
r := &reqWorker{
item: item,
cc: cc,
onClose: onClose,
}
go func(ctx context.Context, r *reqWorker) {
select {
case <-ctx.Done():
if r.onClose != nil {
r.onClose()
}
return
}
}(cctx, r)
return r
}
// callback when a request is evicted.
func onReqCacheEvict(key interface{}, value interface{}) {
// cancel context
value.(*reqWorker).cc()
}
// newRequestWorkerPool .
func newRequestWorkerPool(size int) (*requestWorkerPool, error) {
l, err := lru.NewWithEvict(size, onReqCacheEvict)
return &requestWorkerPool{
cache: l,
}, err
}
// AddReqItem .
// Add with the same reqid will make the previous one cancelled.
// When context is done, item will be removed from the cache;
// When the item is evicted, the context will be cancelled.
func (rc *requestWorkerPool) AddReqItem(ctx context.Context, reqid RequestID, item *reqItem) {
rc.RemoveReqItem(reqid)
w := newReqWorker(ctx, item, func() {
rc.RemoveReqItem(reqid)
})
rc.cache.Add(reqid, w)
}
// RemoveReqItem .
func (rc *requestWorkerPool) RemoveReqItem(reqid RequestID) {
rc.cache.Remove(reqid)
}
// GetReqItem .
func (rc *requestWorkerPool) GetReqItem(reqid RequestID) (out *reqItem, ok bool, cancel func()) {
var w *reqWorker
w, ok = rc.GetReqWorker(reqid)
if w != nil {
out, cancel = w.item, w.cc
}
return
}
// GetReqItem .
func (rc *requestWorkerPool) GetReqWorker(reqid RequestID) (w *reqWorker, ok bool) {
var v interface{}
v, ok = rc.cache.Get(reqid)
if ok {
w = v.(*reqWorker)
}
return
}
// RemoveTopic .
func (rc *requestWorkerPool) RemoveTopic(topic string) {
for _, k := range rc.cache.Keys() {
if v, ok := rc.cache.Peek(k); ok {
w := v.(*reqWorker)
if w.item.topic == topic {
rc.cache.Remove(k)
}
}
}
}
// RemoveAll .
func (rc *requestWorkerPool) RemoveAll() {
rc.cache.Purge()
}
// ResponseCache .
// ResponseCache is used to deduplicate the response
type responseCache struct {
cache *lru.Cache
}
// respItem .
type respItem struct{}
// newResponseCache .
func newResponseCache(size int) (*responseCache, error) {
l, err := lru.New(size)
return &responseCache{
cache: l,
}, err
}
func (r *responseCache) markSeen(resp *Response) bool {
var (
err error
hash uint64
found bool
)
if err == nil {
s := fnv.New64()
_, err = s.Write(resp.Payload)
_, err = s.Write([]byte(resp.Control.RequestId))
hash = s.Sum64()
}
if err == nil {
_, found = r.cache.Get(hash)
}
if err == nil && !found {
r.cache.Add(hash, respItem{})
return true
}
return false
}