This repository has been archived by the owner on Aug 21, 2021. It is now read-only.
/
topic_consumer.go
87 lines (77 loc) · 2.26 KB
/
topic_consumer.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
package channels
//TODO: Test channels!
import (
"gopkg.in/redis.v3"
// "time"
// "fmt"
)
type topicConsumerChannel struct {
redisClient *redis.Client
pubsub *redis.PubSub
channelName string
consumers []Consumer
consumingStopped bool
}
// NewTopicConsumerChannel returns a ConsumerChannel that uses Redis PubSub for
// communication. Each message delivered through this consumer channel will be
// delivered once to each consumer. Note, however, that network issues that
// prevent delivery of a message may lead to messages going completely
// undelivered. Consumers may Ack or Reject the messages, but this is a no-op.
func NewTopicConsumerChannel(channelName string, redisClient *redis.Client) ConsumerChannel {
return &topicConsumerChannel{
redisClient,
nil,
channelName,
[]Consumer{},
false,
}
}
// ReturnAllUnacked is just here for API Compatibility with topics. It does
// nothing
func (topic *topicConsumerChannel) ReturnAllUnacked() int {
return 0
}
// PurgeRejected is just here for API Compatibility with topics. It does
// nothing
func (topic *topicConsumerChannel) PurgeRejected() int {
return 0
}
func (topic *topicConsumerChannel) AddConsumer(consumer Consumer) bool {
topic.consumers = append(topic.consumers, consumer)
return true
}
func (topic *topicConsumerChannel) StartConsuming() bool {
// log.Printf("rmq topic started consuming %s %d %s", topic, prefetchLimit, pollDuration)
if topic.pubsub != nil {
// Already consuming
return false
}
if pubsub, err := topic.redisClient.Subscribe(topic.channelName); err == nil {
topic.pubsub = pubsub
go topic.consume()
return true
}
return false
}
func (topic *topicConsumerChannel) consume() {
for {
msg, err := topic.pubsub.ReceiveMessage()
if err == nil {
for _, consumer := range topic.consumers {
go consumer.Consume(newTopicDelivery(msg.Payload, topic.redisClient))
}
} else if err.Error() == "redis: client is closed" {
return
}
}
}
func (topic *topicConsumerChannel) StopConsuming() bool {
if topic.pubsub != nil && !topic.consumingStopped {
topic.pubsub.Close()
topic.consumingStopped = true
}
return false
}
func (topic *topicConsumerChannel) Publisher() Publisher {
return NewRedisTopicPublisher(topic.channelName, topic.redisClient)
}