代码拉取完成,页面将自动刷新
package natsx
import (
"encoding/json"
"fmt"
"gitee.com/fmpt/fmgo/logx"
"github.com/nats-io/nats.go"
)
func NewNatsClient(addr string) *NatsClient {
var client NatsClient
client.address = fmt.Sprintf("nats://%s", addr)
client.topics = make(map[string]HFunc)
return &client
}
//重新连接 处理
func (c *NatsClient) onReconnected(conn *nats.Conn) {
c.conn = conn
c.subscribe()
}
//失去连接 处理
func (c *NatsClient) onDisconnected(conn *nats.Conn, err error) {
}
func (c *NatsClient) subscribe() {
for topic, handlerFunc := range c.topics {
// 必须用副本
t, h := topic, handlerFunc
if _, err := c.conn.Subscribe(t, func(msg *nats.Msg) {
h(msg.Subject, msg.Reply, msg.Data)
}); err != nil {
logx.Error("Subscribe Topic error : ", err)
continue
}
}
}
// Subscribe 监听主题 必须在 Start 后
func (c *NatsClient) Subscribe(topic string, h HFunc) {
if _, err := c.conn.Subscribe(topic, func(msg *nats.Msg) {
h(msg.Subject, msg.Reply, msg.Data)
}); err != nil {
logx.Error("Subscribe Topic error : ", err)
}
}
func (c *NatsClient) AddTopic(topic string, hFunc HFunc) {
c.topics[topic] = hFunc
}
func (c *NatsClient) AddTopics(topics map[string]HFunc) {
c.topics = topics
}
func (c *NatsClient) Push(topic string, param *NatsMessage) error {
data, err := json.Marshal(param)
if err != nil {
return err
}
if err = c.conn.Publish(topic, data); err != nil {
return err
}
return nil
}
func (c *NatsClient) Push2(topic string, param interface{}) error {
data, err := json.Marshal(param)
if err != nil {
return err
}
if err = c.conn.Publish(topic, data); err != nil {
return err
}
return nil
}
func (c *NatsClient) Response(topic string, param *NatsMessageRsp) error {
data, err := json.Marshal(param)
if err != nil {
return err
}
if err := c.conn.Publish(topic, data); err != nil {
return err
}
return nil
}
func (c *NatsClient) Request(topic string, param *NatsMessage) (*NatsMessageRsp, error) {
data, err := json.Marshal(param)
if err != nil {
return nil, err
}
msg, err := c.conn.Request(topic, data, natsRequestTimeout)
if err != nil {
return nil, err
}
var res NatsMessageRsp
err = json.Unmarshal(msg.Data, &res)
if err != nil {
return nil, err
}
return &res, nil
}
func (c *NatsClient) Start() {
for {
var err error
c.conn, err = nats.Connect(c.address, nats.MaxReconnects(natsReconnectCount), nats.Timeout(natsConnTimeout), nats.ReconnectWait(natsReconnectWaitTime),
nats.DisconnectErrHandler(c.onDisconnected), nats.ReconnectHandler(c.onReconnected))
if err != nil {
logx.Error("Connect to nats svc error : ", err)
continue
}
go c.subscribe()
break
}
}
此处可能存在不合适展示的内容,页面不予展示。您可通过相关编辑功能自查并修改。
如您确认内容无涉及 不当用语 / 纯广告导流 / 暴力 / 低俗色情 / 侵权 / 盗版 / 虚假 / 无价值内容或违法国家有关法律法规的内容,可点击提交进行申诉,我们将尽快为您处理。