2022-02-27 06:57:13 +00:00
|
|
|
package pubsub
|
|
|
|
|
|
|
|
import (
|
|
|
|
"fmt"
|
2022-02-27 10:10:21 +00:00
|
|
|
"log"
|
2022-02-27 06:57:13 +00:00
|
|
|
"sort"
|
|
|
|
|
|
|
|
telemetrypb "git.wntrmute.dev/kyle/sensenet/proto"
|
|
|
|
"git.wntrmute.dev/kyle/sensenet/topic"
|
|
|
|
"google.golang.org/protobuf/proto"
|
|
|
|
"gopkg.in/zeromq/goczmq.v4"
|
|
|
|
)
|
|
|
|
|
2022-02-27 10:10:21 +00:00
|
|
|
func readTopic(b []byte) string {
|
|
|
|
size := int(b[1])
|
|
|
|
return string(b[2 : size+2])
|
|
|
|
}
|
|
|
|
|
|
|
|
func protoTopic(topic string) string {
|
|
|
|
return fmt.Sprintf("\x0a%c%s", len(topic), topic)
|
|
|
|
}
|
|
|
|
|
2022-02-27 06:57:13 +00:00
|
|
|
type Subscriber struct {
|
|
|
|
addr string
|
|
|
|
publisher string
|
|
|
|
sock *goczmq.Sock
|
|
|
|
topics map[string]bool
|
|
|
|
}
|
|
|
|
|
|
|
|
func NewSubscriber(addr, publisher string, topics ...string) (*Subscriber, error) {
|
|
|
|
sub := &Subscriber{
|
2022-02-27 10:10:21 +00:00
|
|
|
addr: addr,
|
|
|
|
publisher: publisher,
|
|
|
|
sock: goczmq.NewSock(goczmq.Sub),
|
|
|
|
topics: map[string]bool{},
|
2022-02-27 06:57:13 +00:00
|
|
|
}
|
|
|
|
|
|
|
|
err := sub.connect()
|
|
|
|
if err != nil {
|
|
|
|
return nil, err
|
|
|
|
}
|
|
|
|
|
|
|
|
sub.Conflate(1)
|
2022-02-27 10:10:21 +00:00
|
|
|
for _, topicName := range topics {
|
|
|
|
sub.Subscribe(topicName)
|
2022-02-27 06:57:13 +00:00
|
|
|
}
|
|
|
|
|
|
|
|
return sub, nil
|
|
|
|
}
|
|
|
|
|
|
|
|
func (sub *Subscriber) Conflate(n int) {
|
|
|
|
sub.sock.SetConflate(n)
|
|
|
|
}
|
|
|
|
|
|
|
|
func (sub *Subscriber) connect() error {
|
2022-02-27 10:10:21 +00:00
|
|
|
log.Printf("subscriber dialing %s", sub.addr)
|
2022-02-27 06:57:13 +00:00
|
|
|
return sub.sock.Connect(sub.addr)
|
|
|
|
}
|
|
|
|
|
|
|
|
func (sub *Subscriber) Subscribe(topic string) {
|
|
|
|
if sub.topics[topic] {
|
|
|
|
return
|
|
|
|
}
|
|
|
|
|
|
|
|
sub.topics[topic] = true
|
2022-02-27 10:10:21 +00:00
|
|
|
sub.sock.SetSubscribe(protoTopic(topic))
|
2022-02-27 06:57:13 +00:00
|
|
|
}
|
|
|
|
|
|
|
|
func (sub *Subscriber) Topics() []string {
|
|
|
|
var topics = make([]string, 0, len(sub.topics))
|
|
|
|
for k := range sub.topics {
|
|
|
|
topics = append(topics, k)
|
|
|
|
}
|
|
|
|
|
|
|
|
sort.Strings(topics)
|
|
|
|
return topics
|
|
|
|
}
|
|
|
|
|
|
|
|
func (sub *Subscriber) receive() ([]byte, error) {
|
|
|
|
var (
|
|
|
|
data []byte
|
|
|
|
err error
|
|
|
|
packet []byte
|
|
|
|
more = 1
|
|
|
|
)
|
|
|
|
for more > 0 {
|
|
|
|
packet, more, err = sub.sock.RecvFrame()
|
|
|
|
if err != nil {
|
|
|
|
return nil, err
|
|
|
|
}
|
|
|
|
|
|
|
|
data = append(data, packet...)
|
|
|
|
}
|
|
|
|
|
|
|
|
return data, nil
|
|
|
|
}
|
|
|
|
|
|
|
|
func (sub *Subscriber) Receive() (*topic.Packet, error) {
|
|
|
|
data, err := sub.receive()
|
|
|
|
if err != nil {
|
|
|
|
return nil, err
|
|
|
|
}
|
|
|
|
|
|
|
|
pbPacket := &telemetrypb.Packet{}
|
|
|
|
err = proto.Unmarshal(data, pbPacket)
|
|
|
|
if err != nil {
|
|
|
|
return nil, err
|
|
|
|
}
|
|
|
|
|
2022-02-27 10:10:21 +00:00
|
|
|
t := readTopic(data)
|
|
|
|
|
2022-02-27 06:57:13 +00:00
|
|
|
packet := &topic.Packet{
|
2022-02-27 10:10:21 +00:00
|
|
|
Topic: t,
|
2022-02-27 06:57:13 +00:00
|
|
|
Publisher: sub.publisher,
|
|
|
|
Received: int64(pbPacket.Timestamp),
|
|
|
|
Payload: pbPacket.Payload,
|
|
|
|
}
|
|
|
|
|
|
|
|
return packet, nil
|
|
|
|
}
|