You can not select more than 25 topics
Topics must start with a letter or number, can include dashes ('-') and can be up to 35 characters long.
141 lines
2.8 KiB
141 lines
2.8 KiB
package mq |
|
|
|
import ( |
|
"github.com/nats-io/nats.go" |
|
"sonet/pkg/utils/logger" |
|
"sync" |
|
"time" |
|
) |
|
|
|
type NatsProducer struct { |
|
options nats.Options |
|
nc *nats.Conn |
|
} |
|
|
|
func NewNatsProducer(options nats.Options) (*NatsProducer, error) { |
|
nc, err := options.Connect() |
|
if err != nil { |
|
return nil, err |
|
} |
|
return &NatsProducer{ |
|
options: options, |
|
nc: nc, |
|
}, nil |
|
} |
|
|
|
func (mq *NatsProducer) Publish(topic string, msg []byte) error { |
|
err := mq.nc.Publish(topic, msg) |
|
if err != nil { |
|
return err |
|
} |
|
return mq.nc.Flush() |
|
} |
|
|
|
func (mq *NatsProducer) MultiPublish(topic string, msgs [][]byte) error { |
|
for _, msg := range msgs { |
|
err := mq.nc.Publish(topic, msg) |
|
if err != nil { |
|
return err |
|
} |
|
} |
|
return mq.nc.Flush() |
|
} |
|
|
|
func (mq *NatsProducer) DelayPublish(topic string, delay time.Duration, body []byte) error { |
|
panic("nats nonsupport delay publish") |
|
} |
|
|
|
func (mq *NatsProducer) Stop() { |
|
err := mq.nc.Flush() |
|
if err != nil { |
|
logger.Error("flush nats error: ", err) |
|
} |
|
mq.nc.Close() |
|
} |
|
|
|
type NatsConsumer struct { |
|
options nats.Options |
|
nc *nats.Conn |
|
consumers map[string]*nats.Subscription |
|
lock *sync.Mutex |
|
} |
|
|
|
func NewNatsConsumer(options nats.Options) (*NatsConsumer, error) { |
|
nc, err := options.Connect() |
|
if err != nil { |
|
return nil, err |
|
} |
|
return &NatsConsumer{ |
|
options: options, |
|
nc: nc, |
|
consumers: make(map[string]*nats.Subscription), |
|
lock: &sync.Mutex{}, |
|
}, nil |
|
} |
|
|
|
func (mq *NatsConsumer) Subscribe(topic string, channel string, handler func(*Message) error) error { |
|
mq.lock.Lock() |
|
defer mq.lock.Unlock() |
|
|
|
if err := mq.unsubscribe(topic); err != nil { |
|
return nil |
|
} |
|
|
|
sub, err := mq.nc.QueueSubscribe(topic, channel, func(msg *nats.Msg) { |
|
message := &Message{Body: msg.Data, Time: time.Now().UnixMilli()} |
|
err := handler(message) |
|
if err != nil { |
|
logger.Error("nats subscribe handle error: ", err) |
|
} |
|
}) |
|
|
|
if err != nil { |
|
return err |
|
} |
|
mq.consumers[topic] = sub |
|
return nil |
|
} |
|
|
|
func (mq *NatsConsumer) SubscribeBroadcast(topic string, handler func(*Message) error) error { |
|
mq.lock.Lock() |
|
defer mq.lock.Unlock() |
|
|
|
if err := mq.unsubscribe(topic); err != nil { |
|
return nil |
|
} |
|
|
|
sub, err := mq.nc.Subscribe(topic, func(msg *nats.Msg) { |
|
message := &Message{Body: msg.Data, Time: time.Now().UnixMilli()} |
|
err := handler(message) |
|
if err != nil { |
|
logger.Error("nats subscribe handle error: ", err) |
|
} |
|
}) |
|
|
|
if err != nil { |
|
return err |
|
} |
|
mq.consumers[topic] = sub |
|
return nil |
|
} |
|
|
|
func (mq *NatsConsumer) Stop() { |
|
mq.lock.Lock() |
|
defer mq.lock.Unlock() |
|
|
|
for topic, sub := range mq.consumers { |
|
err := sub.Unsubscribe() |
|
if err != nil { |
|
logger.Errorf("nats unsubscribe error: topic=%s, err=%v\n", topic, err) |
|
} |
|
} |
|
} |
|
|
|
func (mq *NatsConsumer) unsubscribe(topic string) error { |
|
oldSub := mq.consumers[topic] |
|
if oldSub != nil { |
|
delete(mq.consumers, topic) |
|
return oldSub.Unsubscribe() |
|
} |
|
return nil |
|
}
|
|
|