2024-07-19 20:02:18 +08:00
|
|
|
package nats
|
|
|
|
|
|
2024-07-21 23:09:29 +08:00
|
|
|
import (
|
|
|
|
|
"caatsm/internal/config"
|
|
|
|
|
"context"
|
2024-08-12 17:47:34 +08:00
|
|
|
"errors"
|
2024-07-21 23:09:29 +08:00
|
|
|
|
|
|
|
|
"github.com/ThreeDotsLabs/watermill"
|
|
|
|
|
"github.com/ThreeDotsLabs/watermill-nats/v2/pkg/nats"
|
2024-08-12 17:47:34 +08:00
|
|
|
"github.com/ThreeDotsLabs/watermill/message"
|
2024-07-21 23:09:29 +08:00
|
|
|
nc "github.com/nats-io/nats.go"
|
|
|
|
|
)
|
|
|
|
|
|
2024-08-14 09:40:47 +08:00
|
|
|
func Subscribe(config *config.Config) {
|
2024-07-21 23:09:29 +08:00
|
|
|
logger := watermill.NewStdLogger(false, false)
|
2024-08-14 09:40:47 +08:00
|
|
|
marshaler := &PlainTextMarshaler{}
|
2024-07-21 23:09:29 +08:00
|
|
|
options := []nc.Option{
|
|
|
|
|
nc.RetryOnFailedConnect(true),
|
2024-07-30 10:25:59 +08:00
|
|
|
nc.Timeout(config.Timeouts.Server),
|
|
|
|
|
nc.ReconnectWait(config.Timeouts.ReconnectWait),
|
2024-07-21 23:09:29 +08:00
|
|
|
}
|
|
|
|
|
jsConfig := nats.JetStreamConfig{Disabled: true}
|
|
|
|
|
|
|
|
|
|
subscriber, err := nats.NewSubscriber(
|
|
|
|
|
nats.SubscriberConfig{
|
2024-07-30 10:25:59 +08:00
|
|
|
URL: config.Nats.URL,
|
|
|
|
|
CloseTimeout: config.Timeouts.Close,
|
|
|
|
|
AckWaitTimeout: config.Timeouts.AckWait,
|
2024-07-21 23:09:29 +08:00
|
|
|
NatsOptions: options,
|
|
|
|
|
Unmarshaler: marshaler,
|
|
|
|
|
JetStream: jsConfig,
|
|
|
|
|
},
|
|
|
|
|
logger,
|
|
|
|
|
)
|
|
|
|
|
if err != nil {
|
|
|
|
|
panic(err)
|
|
|
|
|
}
|
2024-08-05 15:58:46 +08:00
|
|
|
|
|
|
|
|
logger.Info("NATS server connected", map[string]interface{}{"url": config.Nats.URL})
|
2024-07-30 10:25:59 +08:00
|
|
|
logger.Info("Subscribing to NATS topic", map[string]interface{}{"topic": config.Subscription.Topic})
|
2024-07-21 23:09:29 +08:00
|
|
|
|
|
|
|
|
defer subscriber.Close()
|
2024-07-30 10:25:59 +08:00
|
|
|
messages, err := subscriber.Subscribe(context.Background(), config.Subscription.Topic)
|
2024-07-21 23:09:29 +08:00
|
|
|
if err != nil {
|
2024-07-30 10:25:59 +08:00
|
|
|
logger.Error("Failed to subscribe to NATS topic", err, map[string]interface{}{"topic": config.Subscription.Topic})
|
2024-07-21 23:09:29 +08:00
|
|
|
return
|
|
|
|
|
}
|
2024-08-05 15:58:46 +08:00
|
|
|
|
2024-08-12 17:47:34 +08:00
|
|
|
handlers := New(config)
|
2024-07-21 23:09:29 +08:00
|
|
|
for msg := range messages {
|
2024-08-12 17:47:34 +08:00
|
|
|
if err := handlers.HandleMessage(msg); err == nil {
|
2024-08-07 11:10:08 +08:00
|
|
|
msg.Ack()
|
|
|
|
|
} else {
|
|
|
|
|
logger.Error("Failed to handle message", err, map[string]interface{}{"message": msg})
|
|
|
|
|
msg.Nack()
|
|
|
|
|
}
|
2024-07-21 23:09:29 +08:00
|
|
|
}
|
|
|
|
|
}
|
2024-08-12 17:47:34 +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
|
|
|
|
|
}
|