/
queue.go
68 lines (56 loc) · 1.4 KB
/
queue.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
package redis
import (
"context"
"time"
"github.com/bantublockchain/push-notification-service/internal/queue"
"github.com/gomodule/redigo/redis"
"gitlab.com/pennersr/redq"
)
type redisQueueFactory struct {
pool *redis.Pool
}
type redisQueue struct {
q *redq.RedQueue
}
// NewQueueFactory ...
func NewQueueFactory(url, pwd string) queue.QueueFactory {
qf := &redisQueueFactory{
pool: &redis.Pool{
MaxIdle: 3,
IdleTimeout: 240 * time.Second,
Dial: func() (redis.Conn, error) {
return redis.DialURL(url, redis.DialPassword(pwd))
},
},
}
return qf
}
func (rq redisQueue) Queue(msg []byte) (err error) {
return rq.q.Queue(msg)
}
func (rq redisQueue) Get(ctx context.Context) (qm queue.QueuedMessage, err error) {
qm, err = rq.q.Get(ctx)
return
}
func (rq redisQueue) Remove(qm queue.QueuedMessage) (err error) {
return rq.q.Remove(qm.(redq.QueuedMessage))
}
func (rq redisQueue) Requeue(qm queue.QueuedMessage) (err error) {
return rq.q.Requeue(qm.(redq.QueuedMessage))
}
func (rq redisQueue) Shutdown() (err error) {
return rq.q.Close()
}
func (rqf *redisQueueFactory) NewQueue(id string) (q queue.Queue, err error) {
waitingList := ListName(id)
rq, err := redq.NewQueue(rqf.pool, waitingList)
if err != nil {
return
}
q = redisQueue{q: rq}
return
}
// ListName returns the Redis list name used for queueing.
func ListName(serviceID string) string {
return "shove:" + serviceID
}