66 lines
1.4 KiB
Go
66 lines
1.4 KiB
Go
// Package leaf starts an in-process nats-server and optional hub leaf.
|
|
package leaf
|
|
|
|
import (
|
|
"errors"
|
|
"net/url"
|
|
"time"
|
|
|
|
natsserver "github.com/nats-io/nats-server/v2/server"
|
|
nats "github.com/nats-io/nats.go"
|
|
)
|
|
|
|
// Node is an in-process broker plus client.
|
|
type Node struct {
|
|
ns *natsserver.Server
|
|
nc *nats.Conn
|
|
}
|
|
|
|
// Start binds 127.0.0.1:0 and optionally leaf-connects to hub.
|
|
func Start(hub string) (*Node, error) {
|
|
opts := &natsserver.Options{Host: "127.0.0.1", Port: -1}
|
|
if hub != "" {
|
|
u, err := url.Parse(hub)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
opts.LeafNode.Remotes = []*natsserver.RemoteLeafOpts{{URLs: []*url.URL{u}}}
|
|
}
|
|
ns, err := natsserver.NewServer(opts)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
go ns.Start()
|
|
if !ns.ReadyForConnections(5 * time.Second) {
|
|
return nil, errors.New("nats not ready")
|
|
}
|
|
nc, err := nats.Connect(ns.ClientURL())
|
|
if err != nil {
|
|
ns.Shutdown()
|
|
return nil, err
|
|
}
|
|
n := &Node{ns: ns, nc: nc}
|
|
if _, err := n.nc.Subscribe("verae.sm.send", func(msg *nats.Msg) {
|
|
if msg.Reply != "" {
|
|
_ = n.nc.Publish(msg.Reply, []byte(`{"accepted":true}`))
|
|
}
|
|
}); err != nil {
|
|
n.Shutdown()
|
|
return nil, err
|
|
}
|
|
return n, nil
|
|
}
|
|
|
|
// Addr is the in-process client URL.
|
|
func (n *Node) Addr() string { return n.ns.ClientURL() }
|
|
|
|
// Shutdown stops client and server.
|
|
func (n *Node) Shutdown() error {
|
|
if n.nc != nil {
|
|
n.nc.Close()
|
|
}
|
|
if n.ns != nil {
|
|
n.ns.Shutdown()
|
|
}
|
|
return nil
|
|
}
|