refactor: Update nats subscription handling

The code changes in `sub.go` update the handling of NATS subscriptions. The `Subscribe` function now takes in a `config` parameter, allowing for more flexibility in configuring the subscription. Additionally, the `Subscribe` function now uses the `handler` and `marshaler` parameters to handle incoming messages and marshal/unmarshal data, respectively. These changes improve the modularity and extensibility of the code when working with NATS subscriptions.
This commit is contained in:
windyboy
2024-07-30 10:25:59 +08:00
parent ebcca5c1ec
commit 766925f7f9
3 changed files with 86 additions and 74 deletions
+4 -2
View File
@@ -2,6 +2,7 @@ package main
import (
"caatsm/internal/config"
"caatsm/internal/handlers"
"caatsm/internal/nats"
"caatsm/pkg/utils"
@@ -24,6 +25,7 @@ func main() {
fmt.Println("Loaded configuration successfully")
log.Info("Starting nats subscriber")
subscribe := nats.NewNatsHandler(cfg)
subscribe.Subscribe()
handler := handlers.NewNatsHandler(cfg)
nats.Subscribe(cfg, &handlers.PlainTextMarshaler{}, handler)
}
+70
View File
@@ -0,0 +1,70 @@
package handlers
import (
"caatsm/internal/config"
"caatsm/internal/domain"
"caatsm/internal/parsers"
"caatsm/internal/repository"
"caatsm/pkg/utils"
"errors"
"fmt"
"sync"
"github.com/ThreeDotsLabs/watermill"
"github.com/ThreeDotsLabs/watermill/message"
nc "github.com/nats-io/nats.go"
)
type NatsHandler struct {
mu sync.Mutex
config *config.Config
hasuraRepo *repository.HasuraRepository
}
func NewNatsHandler(config *config.Config) *NatsHandler {
return &NatsHandler{config: config, hasuraRepo: repository.NewHasuraRepo(config.Hasura.Endpoint, config.Hasura.Secret)}
}
func (n *NatsHandler) HandleMessage(msg *message.Message) error {
n.mu.Lock()
defer n.mu.Unlock()
log := utils.GetSugaredLogger()
if msg.Payload == nil {
log.Error("empty message")
return fmt.Errorf("empty message")
}
payload := string(msg.Payload)
var parsed *domain.ParsedMessage
if parsed = parsers.Parse(payload); !parsed.Parsed {
log.Infof("not parsed: [%s] : {%s} \n", msg.UUID, msg.Payload)
} else {
log.Infof("parsed [%s]: %v\n", msg.UUID, parsed)
}
n.SaveMessage(parsed, msg.UUID)
return nil
}
func (n *NatsHandler) SaveMessage(parsed *domain.ParsedMessage, uuid string) {
logger := watermill.NewStdLogger(false, false)
if parsed != nil {
parsed.Uuid = uuid
if err := n.hasuraRepo.CreateNew(parsed); err != nil {
logger.Error("error inserting message", err, map[string]interface{}{"message": parsed})
}
}
}
type PlainTextMarshaler struct{}
func (m *PlainTextMarshaler) Marshal(topic string, msg nc.Msg) ([]byte, error) {
return msg.Data, nil
}
func (m *PlainTextMarshaler) Unmarshal(newMsg *nc.Msg) (*message.Message, error) {
if newMsg == nil {
return nil, errors.New("empty message")
}
msg := message.NewMessage(watermill.NewUUID(), newMsg.Data)
return msg, nil
}
+12 -72
View File
@@ -2,46 +2,29 @@ package nats
import (
"caatsm/internal/config"
"caatsm/internal/domain"
"caatsm/internal/parsers"
"caatsm/internal/repository"
"caatsm/pkg/utils"
"caatsm/internal/handlers"
"context"
"errors"
"fmt"
"sync"
"github.com/ThreeDotsLabs/watermill"
"github.com/ThreeDotsLabs/watermill-nats/v2/pkg/nats"
"github.com/ThreeDotsLabs/watermill/message"
nc "github.com/nats-io/nats.go"
)
type NatsHandler struct {
mu sync.Mutex
config *config.Config
hasuraRepo *repository.HasuraRepository
}
func NewNatsHandler(config *config.Config) *NatsHandler {
return &NatsHandler{config: config, hasuraRepo: repository.NewHasuraRepo(config.Hasura.Endpoint, config.Hasura.Secret)}
}
func (n *NatsHandler) Subscribe() {
marshaler := &PlainTextMarshaler{}
func Subscribe(config *config.Config, marshaler *handlers.PlainTextMarshaler, handler *handlers.NatsHandler) {
// marshaler := &PlainTextMarshaler{}
logger := watermill.NewStdLogger(false, false)
options := []nc.Option{
nc.RetryOnFailedConnect(true),
nc.Timeout(n.config.Timeouts.Server),
nc.ReconnectWait(n.config.Timeouts.ReconnectWait),
nc.Timeout(config.Timeouts.Server),
nc.ReconnectWait(config.Timeouts.ReconnectWait),
}
jsConfig := nats.JetStreamConfig{Disabled: true}
subscriber, err := nats.NewSubscriber(
nats.SubscriberConfig{
URL: n.config.Nats.URL,
CloseTimeout: n.config.Timeouts.Close,
AckWaitTimeout: n.config.Timeouts.AckWait,
URL: config.Nats.URL,
CloseTimeout: config.Timeouts.Close,
AckWaitTimeout: config.Timeouts.AckWait,
NatsOptions: options,
Unmarshaler: marshaler,
JetStream: jsConfig,
@@ -51,59 +34,16 @@ func (n *NatsHandler) Subscribe() {
if err != nil {
panic(err)
}
logger.Info("Subscribing to NATS topic", map[string]interface{}{"topic": n.config.Subscription.Topic})
logger.Info("Subscribing to NATS topic", map[string]interface{}{"topic": config.Subscription.Topic})
defer subscriber.Close()
messages, err := subscriber.Subscribe(context.Background(), n.config.Subscription.Topic)
messages, err := subscriber.Subscribe(context.Background(), config.Subscription.Topic)
if err != nil {
logger.Error("Failed to subscribe to NATS topic", err, map[string]interface{}{"topic": n.config.Subscription.Topic})
logger.Error("Failed to subscribe to NATS topic", err, map[string]interface{}{"topic": config.Subscription.Topic})
return
}
for msg := range messages {
n.handleMessage(msg)
handler.HandleMessage(msg)
msg.Ack()
}
}
func (n *NatsHandler) handleMessage(msg *message.Message) error {
n.mu.Lock()
defer n.mu.Unlock()
log := utils.GetSugaredLogger()
if msg.Payload == nil {
log.Error("empty message")
return fmt.Errorf("empty message")
}
payload := string(msg.Payload)
var parsed *domain.ParsedMessage
if parsed = parsers.Parse(payload); !parsed.Parsed {
log.Infof("not parsed: [%s] : {%s} \n", msg.UUID, msg.Payload)
} else {
log.Infof("parsed [%s]: %v\n", msg.UUID, parsed)
}
n.SaveMessage(parsed, msg.UUID)
return nil
}
func (n *NatsHandler) SaveMessage(parsed *domain.ParsedMessage, uuid string) {
logger := watermill.NewStdLogger(false, false)
if parsed != nil {
parsed.Uuid = uuid
if err := n.hasuraRepo.CreateNew(parsed); err != nil {
logger.Error("error inserting message", err, map[string]interface{}{"message": parsed})
}
}
}
type PlainTextMarshaler struct{}
func (m *PlainTextMarshaler) Marshal(topic string, msg nc.Msg) ([]byte, error) {
return msg.Data, nil
}
func (m *PlainTextMarshaler) Unmarshal(newMsg *nc.Msg) (*message.Message, error) {
if newMsg == nil {
return nil, errors.New("empty message")
}
msg := message.NewMessage(watermill.NewUUID(), newMsg.Data)
return msg, nil
}