forked from Jesse0Michael/rmq
-
Notifications
You must be signed in to change notification settings - Fork 0
/
Copy pathdelivery.go
126 lines (102 loc) · 2.65 KB
/
delivery.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
package rmq
import (
"context"
"fmt"
"time"
)
type Delivery interface {
Payload() string
Ack() error
Reject() error
Push() error
}
type redisDelivery struct {
ctx context.Context
payload string
unackedKey string
rejectedKey string
pushKey string
redisClient RedisClient
errChan chan<- error
}
func newDelivery(
ctx context.Context,
payload string,
unackedKey string,
rejectedKey string,
pushKey string,
redisClient RedisClient,
errChan chan<- error,
) *redisDelivery {
return &redisDelivery{
ctx: ctx,
payload: payload,
unackedKey: unackedKey,
rejectedKey: rejectedKey,
pushKey: pushKey,
redisClient: redisClient,
errChan: errChan,
}
}
func (delivery *redisDelivery) String() string {
return fmt.Sprintf("[%s %s]", delivery.payload, delivery.unackedKey)
}
func (delivery *redisDelivery) Payload() string {
return delivery.payload
}
// blocking versions of the functions below with the following behavior:
// 1. return immediately if the operation succeeded or failed with ErrorNotFound
// 2. in case of other redis errors, send them to the errors chan and retry after a sleep
// 3. if redis errors occur after StopConsuming() has been called, ErrorConsumingStopped will be returned
func (delivery *redisDelivery) Ack() error {
errorCount := 0
for {
count, err := delivery.redisClient.LRem(delivery.unackedKey, 1, delivery.payload)
if err == nil { // no redis error
if count == 0 {
return ErrorNotFound
}
return nil
}
// redis error
errorCount++
select { // try to add error to channel, but don't block
case delivery.errChan <- &DeliveryError{Delivery: delivery, RedisErr: err, Count: errorCount}:
default:
}
if err := delivery.ctx.Err(); err != nil {
return ErrorConsumingStopped
}
time.Sleep(time.Second)
}
}
func (delivery *redisDelivery) Reject() error {
return delivery.move(delivery.rejectedKey)
}
func (delivery *redisDelivery) Push() error {
if delivery.pushKey == "" {
return delivery.Reject() // fall back to rejecting
}
return delivery.move(delivery.pushKey)
}
func (delivery *redisDelivery) move(key string) error {
errorCount := 0
for {
_, err := delivery.redisClient.LPush(key, delivery.payload)
if err == nil { // success
break
}
// error
errorCount++
select { // try to add error to channel, but don't block
case delivery.errChan <- &DeliveryError{Delivery: delivery, RedisErr: err, Count: errorCount}:
default:
}
if err := delivery.ctx.Err(); err != nil {
return ErrorConsumingStopped
}
time.Sleep(time.Second)
}
return delivery.Ack()
}
// lower level functions which don't retry but just return the first error