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.
82 lines
1.9 KiB
82 lines
1.9 KiB
package nats
|
|
|
|
import (
|
|
"time"
|
|
|
|
gonats "github.com/nats-io/nats.go"
|
|
)
|
|
|
|
func newSys(options *Options) (sys *natsSys, err error) {
|
|
sys = &natsSys{options: options}
|
|
if err = sys.init(); err != nil {
|
|
return nil, err
|
|
}
|
|
return
|
|
}
|
|
|
|
type natsSys struct {
|
|
options *Options
|
|
conn *gonats.Conn
|
|
js gonats.JetStreamContext
|
|
}
|
|
|
|
func (this *natsSys) init() (err error) {
|
|
this.conn, err = gonats.Connect(this.options.URL,
|
|
gonats.MaxReconnects(this.options.MaxReconnects),
|
|
gonats.ReconnectWait(time.Duration(this.options.ReconnectWait)*time.Second),
|
|
gonats.Name(this.options.Name),
|
|
gonats.DisconnectErrHandler(func(_ *gonats.Conn, e error) {
|
|
this.options.Log.Warnf("sys.nats 连接断开: %v", e)
|
|
}),
|
|
gonats.ReconnectHandler(func(c *gonats.Conn) {
|
|
this.options.Log.Infof("sys.nats 已重连: %s", c.ConnectedUrl())
|
|
}),
|
|
)
|
|
if err != nil {
|
|
return
|
|
}
|
|
if this.js, err = this.conn.JetStream(); err != nil {
|
|
this.conn.Close()
|
|
return
|
|
}
|
|
this.options.Log.Infof("sys.nats 连接成功: %s", this.options.URL)
|
|
return
|
|
}
|
|
|
|
// 以下方法一律做 nil 接收者防护:NewSys 失败时调用方若把 (nil, err) 存进 ISys 接口,
|
|
// 接口不为 nil 但接收者为 nil,读字段会直接 panic。宁可降级返回 nil / ErrNotReady。
|
|
func (this *natsSys) Conn() *gonats.Conn {
|
|
if this == nil {
|
|
return nil
|
|
}
|
|
return this.conn
|
|
}
|
|
|
|
func (this *natsSys) JetStream() gonats.JetStreamContext {
|
|
if this == nil {
|
|
return nil
|
|
}
|
|
return this.js
|
|
}
|
|
|
|
func (this *natsSys) Publish(subject string, data []byte) error {
|
|
if this == nil || this.js == nil {
|
|
return ErrNotReady
|
|
}
|
|
_, err := this.js.Publish(subject, data)
|
|
return err
|
|
}
|
|
|
|
func (this *natsSys) PublishAsync(subject string, data []byte) error {
|
|
if this == nil || this.js == nil {
|
|
return ErrNotReady
|
|
}
|
|
_, err := this.js.PublishAsync(subject, data)
|
|
return err
|
|
}
|
|
|
|
func (this *natsSys) Close() {
|
|
if this != nil && this.conn != nil {
|
|
_ = this.conn.Drain()
|
|
}
|
|
}
|
|
|