Merge pull request #5 from windyboy/introduce-interfaces

introduce-interfaces
This commit is contained in:
windyboy
2024-08-14 11:49:43 +08:00
committed by GitHub
6 changed files with 105 additions and 62 deletions
+6 -2
View File
@@ -3,6 +3,7 @@ package main
import ( import (
"caatsm/internal/config" "caatsm/internal/config"
"caatsm/internal/nats" "caatsm/internal/nats"
"caatsm/internal/repository"
"caatsm/pkg/utils" "caatsm/pkg/utils"
"os" "os"
@@ -82,7 +83,10 @@ func executeListen(c *cli.Context) error {
fmt.Println("Loaded configuration successfully") fmt.Println("Loaded configuration successfully")
log := utils.GetLogger() log := utils.GetLogger()
log.Info("Starting nats subscriber") log.Info("Starting nats subscriber")
// handler := handlers.New(cfg) publisher := nats.NewPub(cfg)
nats.Subscribe(cfg) repository := repository.NewHasura(cfg)
handler := nats.NewHandler(cfg, publisher, repository)
subscriber := nats.NewSub(cfg)
subscriber.Subscribe(cfg, handler)
return nil return nil
} }
+22
View File
@@ -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
View File
@@ -3,62 +3,61 @@ package nats
import ( import (
"caatsm/internal/config" "caatsm/internal/config"
"caatsm/internal/domain" "caatsm/internal/domain"
"caatsm/internal/iface"
"caatsm/internal/parsers" "caatsm/internal/parsers"
"caatsm/internal/repository"
"caatsm/pkg/utils" "caatsm/pkg/utils"
"fmt" "fmt"
"sync" "sync"
"github.com/ThreeDotsLabs/watermill"
"github.com/ThreeDotsLabs/watermill/message"
) )
type MessageHandler struct { type MessageHandler struct {
mu sync.Mutex mu sync.Mutex
config *config.Config 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{ return &MessageHandler{
config: config, 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() handler.mu.Lock()
defer handler.mu.Unlock() defer handler.mu.Unlock()
log := utils.GetSugaredLogger() log := utils.GetSugaredLogger()
if msg.Payload == nil { if msg == nil {
log.Error("empty message") log.Error("empty message")
return fmt.Errorf("empty message") return fmt.Errorf("empty message")
} }
payload := string(msg.Payload) payload := string(msg)
var parsed *domain.ParsedMessage var parsed *domain.ParsedMessage
if parsed = parsers.Parse(payload); !parsed.Parsed { 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 { } 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.repository.CreateNew(parsed, uuid)
handler.Publish(parsed) handler.publisher.Publish(parsed)
return nil return nil
} }
func (n *MessageHandler) SaveMessage(parsed *domain.ParsedMessage, uuid string) { // func (n *MessageHandler) SaveMessage(parsed *domain.ParsedMessage, uuid string) {
logger := watermill.NewStdLogger(false, false) // logger := watermill.NewStdLogger(false, false)
if parsed != nil { // if parsed != nil {
parsed.Uuid = uuid // parsed.Uuid = uuid
if err := n.hasuraRepo.CreateNew(parsed); err != nil { // if err := n.hasuraRepo.CreateNew(parsed); err != nil {
logger.Error("error inserting message", err, map[string]interface{}{"message": parsed}) // logger.Error("error inserting message", err, map[string]interface{}{"message": parsed})
} // }
logger.Info("message inserted", map[string]interface{}{"message": parsed.Uuid}) // logger.Info("message inserted", map[string]interface{}{"message": parsed.Uuid})
} // }
} // }
func (n *MessageHandler) Publish(parsed *domain.ParsedMessage) { // func (n *MessageHandler) Publish(parsed *domain.ParsedMessage) {
if err := Publish(n.config, parsed); err != nil { // if err := Publish(n.config, parsed); err != nil {
utils.GetSugaredLogger().Error("error publishing message", err, map[string]interface{}{"message": parsed}) // utils.GetSugaredLogger().Error("error publishing message", err, map[string]interface{}{"message": parsed})
} // }
} // }
+20 -14
View File
@@ -2,6 +2,7 @@ package nats
import ( import (
"caatsm/internal/config" "caatsm/internal/config"
"caatsm/pkg/utils"
"encoding/json" "encoding/json"
"github.com/ThreeDotsLabs/watermill" "github.com/ThreeDotsLabs/watermill"
@@ -10,38 +11,43 @@ import (
nc "github.com/nats-io/nats.go" 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) logger := watermill.NewStdLogger(false, false)
jsConfig := nats.JetStreamConfig{Disabled: true}
options := []nc.Option{ options := []nc.Option{
nc.RetryOnFailedConnect(true), nc.RetryOnFailedConnect(true),
nc.Timeout(config.Timeouts.Server), nc.Timeout(config.Timeouts.Server),
nc.ReconnectWait(config.Timeouts.ReconnectWait), nc.ReconnectWait(config.Timeouts.ReconnectWait),
} }
jsConfig := nats.JetStreamConfig{Disabled: true} publisher, _ := nats.NewPublisher(
publisher, err := nats.NewPublisher(
nats.PublisherConfig{ nats.PublisherConfig{
URL: config.Nats.URL, URL: config.Nats.URL,
NatsOptions: options, NatsOptions: options,
JetStream: jsConfig, JetStream: jsConfig,
}, }, logger)
logger, return &NatsPublisher{
) config: config,
if err != nil { publisher: publisher,
panic(err)
} }
}
logger.Info("NATS server connected", map[string]interface{}{"url": config.Nats.URL}) func (n *NatsPublisher) Publish(parsedMessage interface{}) error {
logger.Info("Publishing message to NATS topic", map[string]interface{}{"topic": config.Publisher.Topic}) logger := utils.GetSugaredLogger()
messageText, err := json.Marshal(parsedMessage) messageText, err := json.Marshal(parsedMessage)
if err != nil { 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)) 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 { 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 err
} }
return nil return nil
+19 -12
View File
@@ -2,6 +2,8 @@ package nats
import ( import (
"caatsm/internal/config" "caatsm/internal/config"
"caatsm/internal/iface"
"caatsm/pkg/utils"
"context" "context"
"errors" "errors"
@@ -11,7 +13,12 @@ import (
nc "github.com/nats-io/nats.go" 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) logger := watermill.NewStdLogger(false, false)
marshaler := &PlainTextMarshaler{} marshaler := &PlainTextMarshaler{}
options := []nc.Option{ options := []nc.Option{
@@ -20,8 +27,7 @@ func Subscribe(config *config.Config) {
nc.ReconnectWait(config.Timeouts.ReconnectWait), nc.ReconnectWait(config.Timeouts.ReconnectWait),
} }
jsConfig := nats.JetStreamConfig{Disabled: true} jsConfig := nats.JetStreamConfig{Disabled: true}
subscriber, _ := nats.NewSubscriber(
subscriber, err := nats.NewSubscriber(
nats.SubscriberConfig{ nats.SubscriberConfig{
URL: config.Nats.URL, URL: config.Nats.URL,
CloseTimeout: config.Timeouts.Close, CloseTimeout: config.Timeouts.Close,
@@ -32,23 +38,24 @@ func Subscribe(config *config.Config) {
}, },
logger, logger,
) )
if err != nil { return &NatsSubscriber{
panic(err) config: config,
subscriber: subscriber,
} }
}
logger.Info("NATS server connected", map[string]interface{}{"url": config.Nats.URL}) func (n *NatsSubscriber) Subscribe(config *config.Config, handlers iface.MessageHandler) {
logger.Info("Subscribing to NATS topic", map[string]interface{}{"topic": config.Subscription.Topic}) logger := utils.GetSugaredLogger()
defer subscriber.Close() defer n.subscriber.Close()
messages, err := subscriber.Subscribe(context.Background(), config.Subscription.Topic) messages, err := n.subscriber.Subscribe(context.Background(), config.Subscription.Topic)
if err != nil { 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 return
} }
handlers := New(config)
for msg := range messages { for msg := range messages {
if err := handlers.HandleMessage(msg); err == nil { if err := handlers.HandleMessage(msg.Payload, []byte(msg.UUID)); err == nil {
msg.Ack() msg.Ack()
} else { } else {
logger.Error("Failed to handle message", err, map[string]interface{}{"message": msg}) logger.Error("Failed to handle message", err, map[string]interface{}{"message": msg})
+10 -5
View File
@@ -5,6 +5,7 @@ import (
"encoding/json" "encoding/json"
"os" "os"
"caatsm/internal/config"
"caatsm/internal/domain" "caatsm/internal/domain"
"caatsm/pkg/utils" "caatsm/pkg/utils"
@@ -18,18 +19,22 @@ type HasuraRepository struct {
} }
// New creates a new HasuraRepository // 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( src := oauth2.StaticTokenSource(
&oauth2.Token{AccessToken: os.Getenv("GRAPHQL_TOKEN")}, &oauth2.Token{AccessToken: token},
) )
httpClient := oauth2.NewClient(context.Background(), src) httpClient := oauth2.NewClient(context.Background(), src)
return &HasuraRepository{ return &HasuraRepository{
client: graphql.NewClient(endpoint, httpClient), client: graphql.NewClient(config.Hasura.Endpoint, httpClient),
} }
} }
// InsertParsedMessage inserts a new ParsedMessage into the Hasura GraphQL API // 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() log := utils.GetSugaredLogger()
bodyString, _ := json.Marshal(pm.BodyData) bodyString, _ := json.Marshal(pm.BodyData)
secondAddress, _ := json.Marshal(pm.SecondaryAddresses) secondAddress, _ := json.Marshal(pm.SecondaryAddresses)
@@ -43,7 +48,7 @@ func (hr *HasuraRepository) CreateNew(pm *domain.ParsedMessage) error {
Category: pm.Category, Category: pm.Category,
Date_time: pm.DateTime, Date_time: pm.DateTime,
Dispatched_at: pm.DispatchedAt, Dispatched_at: pm.DispatchedAt,
Uuid: uuid.New(), Uuid: uuid.Must(uuid.FromBytes(msg_uuid)),
Received_at: pm.ReceivedAt, Received_at: pm.ReceivedAt,
Originator: pm.Originator, Originator: pm.Originator,
Originator_date_time: pm.OriginatorDateTime, Originator_date_time: pm.OriginatorDateTime,