2024-07-19 20:02:18 +08:00
|
|
|
package nats
|
|
|
|
|
|
2024-07-21 23:09:29 +08:00
|
|
|
import (
|
|
|
|
|
"caatsm/internal/config"
|
2024-07-23 16:45:03 +08:00
|
|
|
"caatsm/internal/domain"
|
2024-07-21 23:09:29 +08:00
|
|
|
"caatsm/internal/parsers"
|
2024-07-22 18:00:51 +08:00
|
|
|
"caatsm/internal/repository"
|
2024-07-24 11:07:04 +08:00
|
|
|
"caatsm/pkg/utils"
|
2024-07-21 23:09:29 +08:00
|
|
|
"context"
|
|
|
|
|
"errors"
|
|
|
|
|
"fmt"
|
2024-07-24 14:21:02 +08:00
|
|
|
"sync"
|
2024-07-21 23:09:29 +08:00
|
|
|
|
|
|
|
|
"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 {
|
2024-07-24 14:21:02 +08:00
|
|
|
mu sync.Mutex
|
2024-07-22 18:00:51 +08:00
|
|
|
config *config.Config
|
|
|
|
|
hasuraRepo *repository.HasuraRepository
|
2024-07-21 23:09:29 +08:00
|
|
|
}
|
|
|
|
|
|
|
|
|
|
func NewNatsHandler(config *config.Config) *NatsHandler {
|
2024-07-22 18:00:51 +08:00
|
|
|
return &NatsHandler{config: config, hasuraRepo: repository.NewHasuraRepo(config.Hasura.Endpoint, config.Hasura.Secret)}
|
2024-07-21 23:09:29 +08:00
|
|
|
}
|
|
|
|
|
|
|
|
|
|
func (n *NatsHandler) Subscribe() {
|
|
|
|
|
marshaler := &PlainTextMarshaler{}
|
|
|
|
|
logger := watermill.NewStdLogger(false, false)
|
|
|
|
|
options := []nc.Option{
|
|
|
|
|
nc.RetryOnFailedConnect(true),
|
2024-07-22 12:42:30 +08:00
|
|
|
nc.Timeout(n.config.Timeouts.Server),
|
2024-07-21 23:09:29 +08:00
|
|
|
nc.ReconnectWait(n.config.Timeouts.ReconnectWait),
|
|
|
|
|
}
|
|
|
|
|
jsConfig := nats.JetStreamConfig{Disabled: true}
|
|
|
|
|
|
|
|
|
|
subscriber, err := nats.NewSubscriber(
|
|
|
|
|
nats.SubscriberConfig{
|
|
|
|
|
URL: n.config.Nats.URL,
|
2024-07-22 12:42:30 +08:00
|
|
|
CloseTimeout: n.config.Timeouts.Close,
|
|
|
|
|
AckWaitTimeout: n.config.Timeouts.AckWait,
|
2024-07-21 23:09:29 +08:00
|
|
|
NatsOptions: options,
|
|
|
|
|
Unmarshaler: marshaler,
|
|
|
|
|
JetStream: jsConfig,
|
|
|
|
|
},
|
|
|
|
|
logger,
|
|
|
|
|
)
|
|
|
|
|
if err != nil {
|
|
|
|
|
panic(err)
|
|
|
|
|
}
|
|
|
|
|
logger.Info("Subscribing to NATS topic", map[string]interface{}{"topic": n.config.Subscription.Topic})
|
|
|
|
|
|
|
|
|
|
defer subscriber.Close()
|
|
|
|
|
messages, err := subscriber.Subscribe(context.Background(), n.config.Subscription.Topic)
|
|
|
|
|
if err != nil {
|
|
|
|
|
logger.Error("Failed to subscribe to NATS topic", err, map[string]interface{}{"topic": n.config.Subscription.Topic})
|
|
|
|
|
return
|
|
|
|
|
}
|
|
|
|
|
for msg := range messages {
|
|
|
|
|
n.handleMessage(msg)
|
|
|
|
|
msg.Ack()
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
func (n *NatsHandler) handleMessage(msg *message.Message) error {
|
2024-07-24 14:21:02 +08:00
|
|
|
n.mu.Lock()
|
|
|
|
|
defer n.mu.Unlock()
|
2024-07-24 11:07:04 +08:00
|
|
|
log := utils.GetSugaredLogger()
|
2024-07-21 23:09:29 +08:00
|
|
|
if msg.Payload == nil {
|
2024-07-24 11:07:04 +08:00
|
|
|
log.Error("empty message")
|
2024-07-21 23:09:29 +08:00
|
|
|
return fmt.Errorf("empty message")
|
|
|
|
|
}
|
|
|
|
|
payload := string(msg.Payload)
|
2024-07-23 16:45:03 +08:00
|
|
|
var parsed *domain.ParsedMessage
|
2024-07-26 12:49:01 +08:00
|
|
|
if parsed = parsers.Parse(payload); !parsed.Parsed {
|
|
|
|
|
log.Infof("not parsed: [%s] : {%s} \n", msg.UUID, msg.Payload)
|
2024-07-21 23:09:29 +08:00
|
|
|
} else {
|
2024-07-24 11:07:04 +08:00
|
|
|
log.Infof("parsed [%s]: %v\n", msg.UUID, parsed)
|
2024-07-23 16:45:03 +08:00
|
|
|
|
2024-07-21 23:09:29 +08:00
|
|
|
}
|
2024-07-26 12:49:01 +08:00
|
|
|
n.SaveMessage(parsed, msg.UUID)
|
|
|
|
|
return nil
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
func (n *NatsHandler) SaveMessage(parsed *domain.ParsedMessage, uuid string) {
|
|
|
|
|
logger := watermill.NewStdLogger(false, false)
|
2024-07-24 14:21:02 +08:00
|
|
|
if parsed != nil {
|
2024-07-26 12:49:01 +08:00
|
|
|
parsed.Uuid = uuid
|
|
|
|
|
if err := n.hasuraRepo.CreateNew(parsed); err != nil {
|
|
|
|
|
logger.Error("error inserting message", err, map[string]interface{}{"message": parsed})
|
2024-07-24 14:21:02 +08:00
|
|
|
}
|
2024-07-23 16:45:03 +08:00
|
|
|
}
|
2024-07-21 23:09:29 +08:00
|
|
|
}
|
|
|
|
|
|
|
|
|
|
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
|
2024-07-19 20:02:18 +08:00
|
|
|
}
|