/
nsq.go
76 lines (59 loc) · 1.29 KB
/
nsq.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
package main
import (
"fmt"
"time"
"github.com/golang/protobuf/proto"
"github.com/golang/protobuf/ptypes"
"github.com/jnewmano/mqtt-nsq/mqtt"
"github.com/jnewmano/mqtt-nsq/nsqexporter"
"github.com/nsqio/go-nsq"
)
type nsqProducer struct {
producer *nsq.Producer
topic string
wrap bool
}
func newNSQProducer(addr string, topic string, wrap bool) (*nsqProducer, error) {
config := nsq.NewConfig()
np, err := nsq.NewProducer(addr, config)
if err != nil {
return nil, err
}
// TODO: intercept log messages
np.SetLogger(nil, nsq.LogLevelDebug)
p := nsqProducer{
producer: np,
topic: topic,
wrap: wrap,
}
return &p, nil
}
func (n *nsqProducer) Handle(p *mqtt.Publish) error {
var body []byte
if n.wrap {
t, err := ptypes.TimestampProto(time.Now())
if err != nil {
return err
}
m := nsqexporter.MQTTMessage{
Timestamp: t,
Topic: p.Topic,
Payload: p.Payload,
SourceAddress: "", // we don't know where the message originated
PacketID: uint32(p.PacketID),
}
body, err = proto.Marshal(&m)
if err != nil {
return err
}
} else {
body = p.Payload
}
err := n.producer.Publish(n.topic, body)
if err != nil {
fmt.Printf("unable to publish message to NSQ [%s]\n", err)
return err
}
fmt.Printf(".")
return nil
}