-
Notifications
You must be signed in to change notification settings - Fork 0
/
qreader.go
55 lines (48 loc) · 1000 Bytes
/
qreader.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
package main
import (
"github.com/streadway/amqp"
"log"
)
type QueueReader struct {
conn *amqp.Connection
channel *amqp.Channel
inbound <-chan amqp.Delivery // receive only channel
}
func NewQueueReader(url string) (*QueueReader, error) {
qreader := QueueReader{}
var err error
qreader.conn, err = amqp.Dial(url)
if err != nil {
log.Printf("Failed to connect to queue address url: %s\n", url)
return nil, err
}
qreader.channel, err = qreader.conn.Channel()
if err != nil {
log.Printf("Failed to create channel\n")
return nil, err
}
return &qreader, nil
}
func (q *QueueReader) Consume(queue string) (<-chan []byte, error) {
var err error
outbound := make(chan []byte)
q.inbound, err = q.channel.Consume(
queue,
"qreader",
true,
false,
false,
false,
nil,
)
if err != nil {
log.Printf("Failed to create consumer on channel\n")
return nil, err
}
go func() {
for msg := range q.inbound {
outbound <- msg.Body
}
}()
return outbound, nil
}