From 766925f7f966b89fca571389b81bc6ed7d5bf471 Mon Sep 17 00:00:00 2001 From: windyboy Date: Tue, 30 Jul 2024 10:25:59 +0800 Subject: [PATCH] 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. --- cmd/main/main.go | 6 ++- internal/handlers/nats_handler.go | 70 ++++++++++++++++++++++++++ internal/nats/sub.go | 84 +++++-------------------------- 3 files changed, 86 insertions(+), 74 deletions(-) create mode 100644 internal/handlers/nats_handler.go diff --git a/cmd/main/main.go b/cmd/main/main.go index 347be07..a0393f9 100644 --- a/cmd/main/main.go +++ b/cmd/main/main.go @@ -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) } diff --git a/internal/handlers/nats_handler.go b/internal/handlers/nats_handler.go new file mode 100644 index 0000000..da002a1 --- /dev/null +++ b/internal/handlers/nats_handler.go @@ -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 +} diff --git a/internal/nats/sub.go b/internal/nats/sub.go index b8c8657..443abde 100644 --- a/internal/nats/sub.go +++ b/internal/nats/sub.go @@ -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 -}