-
Notifications
You must be signed in to change notification settings - Fork 7
/
response_callback.go
82 lines (63 loc) · 1.9 KB
/
response_callback.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
package main
import (
"sync"
"time"
)
// ResponseCallback the response from the thrift
type ResponseCallback = func(response *Message, err error)
type responseWithTimeout struct {
responseCallback ResponseCallback
timeoutTime time.Time
}
func newResponseWithTimeout(responseCallback ResponseCallback,
timeoutTime time.Time) *responseWithTimeout {
return &responseWithTimeout{responseCallback: responseCallback,
timeoutTime: timeoutTime}
}
func (r *responseWithTimeout) isTimeout() bool {
return r.timeoutTime.Before(time.Now())
}
type ResponseCallbackMgr struct {
sync.Mutex
responseCallbacks map[int]*responseWithTimeout
}
// NewResponseCallbackMgr create a ResponseCallbackMgr object
func NewResponseCallbackMgr() *ResponseCallbackMgr {
return &ResponseCallbackMgr{responseCallbacks: make(map[int]*responseWithTimeout)}
}
// Add add a response callback for a seqId
func (r *ResponseCallbackMgr) Add(seqId int, callback ResponseCallback, timeoutTime time.Time) {
r.Lock()
defer r.Unlock()
r.responseCallbacks[seqId] = newResponseWithTimeout(callback, timeoutTime)
}
func (r *ResponseCallbackMgr) Remove(seqId int) (ResponseCallback, bool) {
r.Lock()
defer r.Unlock()
if value, ok := r.responseCallbacks[seqId]; ok {
delete(r.responseCallbacks, seqId)
return value.responseCallback, true
} else {
return nil, false
}
}
func (r *ResponseCallbackMgr) getTimeoutResponses() map[int]*responseWithTimeout {
r.Lock()
defer r.Unlock()
timeoutItems := make(map[int]*responseWithTimeout)
for key, value := range r.responseCallbacks {
if value.isTimeout() {
timeoutItems[key] = value
}
}
for key, _ := range timeoutItems {
delete(r.responseCallbacks, key)
}
return timeoutItems
}
func (r *ResponseCallbackMgr) RemoveTimeout(procFunc func(callback ResponseCallback)) {
timeoutItems := r.getTimeoutResponses()
for _, value := range timeoutItems {
procFunc(value.responseCallback)
}
}