✨ Add configuration to interfaces.
This commit is contained in:
@@ -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
|
||||
}
|
||||
+28
-29
@@ -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})
|
||||
// }
|
||||
// }
|
||||
|
||||
+20
-14
@@ -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
|
||||
|
||||
+19
-12
@@ -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})
|
||||
|
||||
@@ -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,
|
||||
|
||||
Reference in New Issue
Block a user