-
Notifications
You must be signed in to change notification settings - Fork 0
/
Copy pathbridge.go
93 lines (78 loc) · 2.49 KB
/
bridge.go
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
package vnats
import (
"fmt"
"log/slog"
"strings"
natsServer "github.com/nats-io/nats-server/v2/server"
"github.com/nats-io/nats.go"
)
type natsBridge struct {
connection *nats.Conn
jetStreamContext nats.JetStreamContext
logger *slog.Logger
}
func newNATSBridge(servers []string, logger *slog.Logger) (*natsBridge, error) {
nb := &natsBridge{
logger: logger,
}
var err error
url := strings.Join(servers, ",")
nb.connection, err = nats.Connect(url,
nats.DisconnectErrHandler(func(_ *nats.Conn, err error) {
logger.Error("Got disconnected", slog.String("error", err.Error()))
}),
nats.ReconnectHandler(func(nc *nats.Conn) {
logger.Error("Got reconnected to!", slog.String("url", nc.ConnectedUrl()))
}),
nats.ClosedHandler(func(nc *nats.Conn) {
logger.Error("Connection closed", slog.String("error", nc.LastError().Error()))
}))
if err != nil {
return nil, fmt.Errorf("could not make NATS Connection to %s: %w", url, err)
}
nb.jetStreamContext, err = nb.connection.JetStream()
if err != nil {
return nil, err
}
return nb, nil
}
func (b *natsBridge) PublishMsg(msg *nats.Msg, msgID string) error {
_, err := b.jetStreamContext.PublishMsg(msg, nats.MsgId(msgID))
return err
}
func (b *natsBridge) EnsureStreamExists(streamConfig *nats.StreamConfig) error {
if _, err := b.jetStreamContext.StreamInfo(streamConfig.Name); err != nil {
if err != nats.ErrStreamNotFound {
return fmt.Errorf("NATS streamInfo-info could not be fetched: %w", err)
}
b.logger.Info("Stream not found, about to add stream.", slog.String("name", streamConfig.Name))
_, err = b.jetStreamContext.AddStream(streamConfig)
if err != nil {
return fmt.Errorf("streamInfo %s could not be added: %w", streamConfig.Name, err)
}
b.logger.Info("Added new NATS streamInfo", slog.String("name", streamConfig.Name))
}
return nil
}
func (b *natsBridge) Subscribe(subject, consumerName string, mode SubscriptionMode) (*nats.Subscription, error) {
var maxAckPending int
switch mode {
case MultipleSubscribersAllowed:
maxAckPending = natsServer.JsDefaultMaxAckPending
case SingleSubscriberStrictMessageOrder:
maxAckPending = 1
default:
maxAckPending = natsServer.JsDefaultMaxAckPending
}
return b.jetStreamContext.PullSubscribe(subject, consumerName,
nats.AckExplicit(),
nats.MaxAckPending(maxAckPending),
nats.AckWait(defaultAckWait),
)
}
func (b *natsBridge) Servers() []string {
return b.connection.Servers()
}
func (b *natsBridge) Drain() error {
return b.connection.Drain()
}