-
Notifications
You must be signed in to change notification settings - Fork 13
/
Copy pathmailbox.go
117 lines (88 loc) · 2.06 KB
/
mailbox.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
package vega
type MemMailbox struct {
name string
values []*Message
inflight map[MessageId]*Message
watchers []*watchChannel
}
func NewMemMailbox(name string) Mailbox {
return &MemMailbox{name, nil, make(map[MessageId]*Message), nil}
}
func (mm *MemMailbox) Ack(id MessageId) error {
if _, ok := mm.inflight[id]; ok {
delete(mm.inflight, id)
return nil
}
return EUnknownMessage
}
func (mm *MemMailbox) Nack(id MessageId) error {
if c, ok := mm.inflight[id]; ok {
delete(mm.inflight, id)
mm.values = append([]*Message{c}, mm.values...)
return nil
}
return EUnknownMessage
}
func (mm *MemMailbox) Abandon() error {
mm.values = nil
for _, w := range mm.watchers {
w.indicator <- nil
}
return nil
}
func (mm *MemMailbox) Poll() (*Message, error) {
if len(mm.values) > 0 {
val := mm.values[0]
mm.values = mm.values[1:]
if val.MessageId == "" {
val.MessageId = NextMessageID()
}
mm.inflight[val.MessageId] = val
return val, nil
}
return nil, nil
}
func (mm *MemMailbox) Push(value *Message) error {
RETRY:
if len(mm.watchers) > 0 {
watch := mm.watchers[0]
mm.watchers = mm.watchers[1:]
if watch.done != nil {
select {
case <-watch.done:
close(watch.indicator)
goto RETRY
default:
}
}
if value.MessageId == "" {
value.MessageId = NextMessageID()
}
mm.inflight[value.MessageId] = value
watch.indicator <- value
close(watch.indicator)
return nil
}
mm.values = append(mm.values, value)
return nil
}
type watchChannel struct {
indicator chan *Message
done chan struct{}
}
func (mm *MemMailbox) AddWatcher() <-chan *Message {
indicator := make(chan *Message, 1)
mm.watchers = append(mm.watchers, &watchChannel{indicator, nil})
return indicator
}
func (mm *MemMailbox) AddWatcherCancelable(done chan struct{}) <-chan *Message {
indicator := make(chan *Message, 1)
mm.watchers = append(mm.watchers, &watchChannel{indicator, done})
return indicator
}
func (mm *MemMailbox) Stats() *MailboxStats {
return &MailboxStats{
Size: len(mm.values),
InFlight: len(mm.inflight),
}
}