diff --git a/cmd/main/main.go b/cmd/main/main.go index 480b8e0..0428acf 100644 --- a/cmd/main/main.go +++ b/cmd/main/main.go @@ -3,6 +3,7 @@ package main import ( "caatsm/internal/config" "caatsm/internal/nats" + "caatsm/internal/repository" "caatsm/pkg/utils" "os" @@ -82,7 +83,10 @@ func executeListen(c *cli.Context) error { fmt.Println("Loaded configuration successfully") log := utils.GetLogger() log.Info("Starting nats subscriber") - // handler := handlers.New(cfg) - nats.Subscribe(cfg) + publisher := nats.NewPub(cfg) + repository := repository.NewHasura(cfg) + handler := nats.NewHandler(cfg, publisher, repository) + subscriber := nats.NewSub(cfg) + subscriber.Subscribe(cfg, handler) return nil } diff --git a/internal/iface/interface.go b/internal/iface/interface.go new file mode 100644 index 0000000..1ea174b --- /dev/null +++ b/internal/iface/interface.go @@ -0,0 +1,22 @@ +package iface + +import ( + "caatsm/internal/config" + "caatsm/internal/domain" +) + +type MessageHandler interface { + HandleMessage(msg []byte, uuid []byte) error +} + +type MessagePublisher interface { + Publish(message interface{}) error +} + +type MessageSubscriber interface { + Subscribe(config *config.Config) error +} + +type MessageRepository interface { + CreateNew(message *domain.ParsedMessage, uuid []byte) error +} diff --git a/internal/nats/handler.go b/internal/nats/handler.go index 4846808..f977904 100644 --- a/internal/nats/handler.go +++ b/internal/nats/handler.go @@ -3,62 +3,61 @@ package nats import ( "caatsm/internal/config" "caatsm/internal/domain" + "caatsm/internal/iface" "caatsm/internal/parsers" - "caatsm/internal/repository" "caatsm/pkg/utils" "fmt" "sync" - - "github.com/ThreeDotsLabs/watermill" - "github.com/ThreeDotsLabs/watermill/message" ) type MessageHandler struct { mu sync.Mutex config *config.Config - hasuraRepo *repository.HasuraRepository + repository iface.MessageRepository + publisher iface.MessagePublisher } -func New(config *config.Config) *MessageHandler { +func NewHandler(config *config.Config, publisher iface.MessagePublisher, repository iface.MessageRepository) *MessageHandler { return &MessageHandler{ config: config, - hasuraRepo: repository.New(config.Hasura.Endpoint, config.Hasura.Secret), + repository: repository, + publisher: publisher, } } -func (handler *MessageHandler) HandleMessage(msg *message.Message) error { +func (handler *MessageHandler) HandleMessage(msg []byte, uuid []byte) error { handler.mu.Lock() defer handler.mu.Unlock() log := utils.GetSugaredLogger() - if msg.Payload == nil { + if msg == nil { log.Error("empty message") return fmt.Errorf("empty message") } - payload := string(msg.Payload) + payload := string(msg) var parsed *domain.ParsedMessage if parsed = parsers.Parse(payload); !parsed.Parsed { - log.Infof("not parsed: [%s] : {%s} \n", msg.UUID, msg.Payload) + log.Infof("not parsed: [%s] : {%s} \n", uuid, payload) } else { - log.Infof("parsed [%s]: %v\n", msg.UUID, parsed.ToString()) + log.Infof("parsed [%s]: %v\n", uuid, parsed.ToString()) } - handler.SaveMessage(parsed, msg.UUID) - handler.Publish(parsed) + handler.repository.CreateNew(parsed, uuid) + handler.publisher.Publish(parsed) return nil } -func (n *MessageHandler) 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}) - } - logger.Info("message inserted", map[string]interface{}{"message": parsed.Uuid}) - } -} +// func (n *MessageHandler) 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}) +// } +// logger.Info("message inserted", map[string]interface{}{"message": parsed.Uuid}) +// } +// } -func (n *MessageHandler) Publish(parsed *domain.ParsedMessage) { - if err := Publish(n.config, parsed); err != nil { - utils.GetSugaredLogger().Error("error publishing message", err, map[string]interface{}{"message": parsed}) - } -} +// func (n *MessageHandler) Publish(parsed *domain.ParsedMessage) { +// if err := Publish(n.config, parsed); err != nil { +// utils.GetSugaredLogger().Error("error publishing message", err, map[string]interface{}{"message": parsed}) +// } +// } diff --git a/internal/nats/pub.go b/internal/nats/pub.go index e8cc5df..9a591d5 100644 --- a/internal/nats/pub.go +++ b/internal/nats/pub.go @@ -2,6 +2,7 @@ package nats import ( "caatsm/internal/config" + "caatsm/pkg/utils" "encoding/json" "github.com/ThreeDotsLabs/watermill" @@ -10,38 +11,43 @@ import ( nc "github.com/nats-io/nats.go" ) -func Publish(config *config.Config, parsedMessage interface{}) error { +type NatsPublisher struct { + config *config.Config + publisher *nats.Publisher +} + +func NewPub(config *config.Config) *NatsPublisher { logger := watermill.NewStdLogger(false, false) + + jsConfig := nats.JetStreamConfig{Disabled: true} options := []nc.Option{ nc.RetryOnFailedConnect(true), nc.Timeout(config.Timeouts.Server), nc.ReconnectWait(config.Timeouts.ReconnectWait), } - jsConfig := nats.JetStreamConfig{Disabled: true} - - publisher, err := nats.NewPublisher( + publisher, _ := nats.NewPublisher( nats.PublisherConfig{ URL: config.Nats.URL, NatsOptions: options, JetStream: jsConfig, - }, - logger, - ) - if err != nil { - panic(err) + }, logger) + return &NatsPublisher{ + config: config, + publisher: publisher, } +} - logger.Info("NATS server connected", map[string]interface{}{"url": config.Nats.URL}) - logger.Info("Publishing message to NATS topic", map[string]interface{}{"topic": config.Publisher.Topic}) +func (n *NatsPublisher) Publish(parsedMessage interface{}) error { + logger := utils.GetSugaredLogger() messageText, err := json.Marshal(parsedMessage) if err != nil { - logger.Error("Failed to marshal message", err, map[string]interface{}{"message": parsedMessage}) + logger.Errorf("Failed to marshal message: %v", err) } msg := message.NewMessage(watermill.NewUUID(), []byte(messageText)) - err = publisher.Publish(config.Publisher.Topic, msg) + err = n.publisher.Publish(n.config.Publisher.Topic, msg) if err != nil { - logger.Error("Failed to publish message to NATS topic", err, map[string]interface{}{"topic": config.Publisher.Topic}) + logger.Errorf("Failed to publish message: %v", err) return err } return nil diff --git a/internal/nats/sub.go b/internal/nats/sub.go index cb8060f..4dbe7f9 100644 --- a/internal/nats/sub.go +++ b/internal/nats/sub.go @@ -2,6 +2,8 @@ package nats import ( "caatsm/internal/config" + "caatsm/internal/iface" + "caatsm/pkg/utils" "context" "errors" @@ -11,7 +13,12 @@ import ( nc "github.com/nats-io/nats.go" ) -func Subscribe(config *config.Config) { +type NatsSubscriber struct { + config *config.Config + subscriber *nats.Subscriber +} + +func NewSub(config *config.Config) *NatsSubscriber { logger := watermill.NewStdLogger(false, false) marshaler := &PlainTextMarshaler{} options := []nc.Option{ @@ -20,8 +27,7 @@ func Subscribe(config *config.Config) { nc.ReconnectWait(config.Timeouts.ReconnectWait), } jsConfig := nats.JetStreamConfig{Disabled: true} - - subscriber, err := nats.NewSubscriber( + subscriber, _ := nats.NewSubscriber( nats.SubscriberConfig{ URL: config.Nats.URL, CloseTimeout: config.Timeouts.Close, @@ -32,23 +38,24 @@ func Subscribe(config *config.Config) { }, logger, ) - if err != nil { - panic(err) + return &NatsSubscriber{ + config: config, + subscriber: subscriber, } +} - logger.Info("NATS server connected", map[string]interface{}{"url": config.Nats.URL}) - logger.Info("Subscribing to NATS topic", map[string]interface{}{"topic": config.Subscription.Topic}) +func (n *NatsSubscriber) Subscribe(config *config.Config, handlers iface.MessageHandler) { + logger := utils.GetSugaredLogger() - defer subscriber.Close() - messages, err := subscriber.Subscribe(context.Background(), config.Subscription.Topic) + defer n.subscriber.Close() + messages, err := n.subscriber.Subscribe(context.Background(), config.Subscription.Topic) if err != nil { - logger.Error("Failed to subscribe to NATS topic", err, map[string]interface{}{"topic": config.Subscription.Topic}) + logger.Errorf("Failed to subscribe to topic: %v", err) return } - handlers := New(config) for msg := range messages { - if err := handlers.HandleMessage(msg); err == nil { + if err := handlers.HandleMessage(msg.Payload, []byte(msg.UUID)); err == nil { msg.Ack() } else { logger.Error("Failed to handle message", err, map[string]interface{}{"message": msg}) diff --git a/internal/repository/hasura.go b/internal/repository/hasura.go index 289ff32..abe2d2b 100644 --- a/internal/repository/hasura.go +++ b/internal/repository/hasura.go @@ -5,6 +5,7 @@ import ( "encoding/json" "os" + "caatsm/internal/config" "caatsm/internal/domain" "caatsm/pkg/utils" @@ -18,18 +19,22 @@ type HasuraRepository struct { } // New creates a new HasuraRepository -func New(endpoint, secret string) *HasuraRepository { +func NewHasura(config *config.Config) *HasuraRepository { + token := os.Getenv("GRAPHQL_TOKEN") + if token == "" { + token = config.Hasura.Secret + } src := oauth2.StaticTokenSource( - &oauth2.Token{AccessToken: os.Getenv("GRAPHQL_TOKEN")}, + &oauth2.Token{AccessToken: token}, ) httpClient := oauth2.NewClient(context.Background(), src) return &HasuraRepository{ - client: graphql.NewClient(endpoint, httpClient), + client: graphql.NewClient(config.Hasura.Endpoint, httpClient), } } // InsertParsedMessage inserts a new ParsedMessage into the Hasura GraphQL API -func (hr *HasuraRepository) CreateNew(pm *domain.ParsedMessage) error { +func (hr *HasuraRepository) CreateNew(pm *domain.ParsedMessage, msg_uuid []byte) error { log := utils.GetSugaredLogger() bodyString, _ := json.Marshal(pm.BodyData) secondAddress, _ := json.Marshal(pm.SecondaryAddresses) @@ -43,7 +48,7 @@ func (hr *HasuraRepository) CreateNew(pm *domain.ParsedMessage) error { Category: pm.Category, Date_time: pm.DateTime, Dispatched_at: pm.DispatchedAt, - Uuid: uuid.New(), + Uuid: uuid.Must(uuid.FromBytes(msg_uuid)), Received_at: pm.ReceivedAt, Originator: pm.Originator, Originator_date_time: pm.OriginatorDateTime,