From 8a4c46aa34bdbd157f47ef20bafaeaa8084f589b Mon Sep 17 00:00:00 2001 From: windyboy Date: Mon, 12 Aug 2024 17:47:34 +0800 Subject: [PATCH] update handler (#2) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit * 🔧 Remove commented-out regex patterns and initialization in config. * ✨ Add NATS message publisher with configuration 📩🚀 * ✨ Add new PlainTextMarshaler and refactor handler in nats package. --- cmd/main/main.go | 5 +- configs/config.dev.toml | 3 ++ internal/config/config.go | 5 ++ .../message_handler.go => nats/handler.go} | 19 ++------ internal/nats/pub.go | 48 +++++++++++++++++++ internal/nats/sub.go | 22 +++++++-- 6 files changed, 82 insertions(+), 20 deletions(-) rename internal/{handlers/message_handler.go => nats/handler.go} (77%) create mode 100644 internal/nats/pub.go diff --git a/cmd/main/main.go b/cmd/main/main.go index f39f381..d9c13da 100644 --- a/cmd/main/main.go +++ b/cmd/main/main.go @@ -2,7 +2,6 @@ package main import ( "caatsm/internal/config" - "caatsm/internal/handlers" "caatsm/internal/nats" "caatsm/pkg/utils" "os" @@ -83,7 +82,7 @@ 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, &handlers.PlainTextMarshaler{}, handler) + // handler := handlers.New(cfg) + nats.Subscribe(cfg, &nats.PlainTextMarshaler{}) return nil } diff --git a/configs/config.dev.toml b/configs/config.dev.toml index 3bcdbac..be381ba 100644 --- a/configs/config.dev.toml +++ b/configs/config.dev.toml @@ -7,6 +7,9 @@ cluster = "tele-cluster" topic = "Telegram.Serial" queue = "tele-queue" +[publisher] +topic = "Telegram.Json" + [timeouts] server = "5s" reconnect_wait = "5s" diff --git a/internal/config/config.go b/internal/config/config.go index 1593243..062a79f 100644 --- a/internal/config/config.go +++ b/internal/config/config.go @@ -15,6 +15,7 @@ var MyConfig *Config type Config struct { Nats NatsConfig Subscription SubscriptionConfig + Publisher PublisherConfig Timeouts TimeoutsConfig Hasura HasuraConfig } @@ -30,6 +31,10 @@ type SubscriptionConfig struct { QueueGroup string `mapstructure:"queue_group"` } +type PublisherConfig struct { + Topic string `mapstructure:"topic"` +} + type TimeoutsConfig struct { Server time.Duration `mapstructure:"server"` ReconnectWait time.Duration `mapstructure:"reconnect_wait"` diff --git a/internal/handlers/message_handler.go b/internal/nats/handler.go similarity index 77% rename from internal/handlers/message_handler.go rename to internal/nats/handler.go index 62997e1..79b9dd2 100644 --- a/internal/handlers/message_handler.go +++ b/internal/nats/handler.go @@ -1,4 +1,4 @@ -package handlers +package nats import ( "caatsm/internal/config" @@ -6,13 +6,11 @@ import ( "caatsm/internal/parsers" "caatsm/internal/repository" "caatsm/pkg/utils" - "errors" "fmt" "sync" "github.com/ThreeDotsLabs/watermill" "github.com/ThreeDotsLabs/watermill/message" - nc "github.com/nats-io/nats.go" ) type MessageHandler struct { @@ -44,6 +42,7 @@ func (handler *MessageHandler) HandleMessage(msg *message.Message) error { log.Infof("parsed [%s]: %v\n", msg.UUID, parsed.ToString()) } handler.SaveMessage(parsed, msg.UUID) + handler.Publish(parsed) return nil } @@ -57,16 +56,8 @@ func (n *MessageHandler) SaveMessage(parsed *domain.ParsedMessage, uuid string) } } -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") +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}) } - msg := message.NewMessage(watermill.NewUUID(), newMsg.Data) - return msg, nil } diff --git a/internal/nats/pub.go b/internal/nats/pub.go new file mode 100644 index 0000000..e8cc5df --- /dev/null +++ b/internal/nats/pub.go @@ -0,0 +1,48 @@ +package nats + +import ( + "caatsm/internal/config" + "encoding/json" + + "github.com/ThreeDotsLabs/watermill" + "github.com/ThreeDotsLabs/watermill-nats/v2/pkg/nats" + "github.com/ThreeDotsLabs/watermill/message" + nc "github.com/nats-io/nats.go" +) + +func Publish(config *config.Config, parsedMessage interface{}) error { + logger := watermill.NewStdLogger(false, false) + 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( + nats.PublisherConfig{ + URL: config.Nats.URL, + NatsOptions: options, + JetStream: jsConfig, + }, + logger, + ) + if err != nil { + panic(err) + } + + 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}) + + messageText, err := json.Marshal(parsedMessage) + if err != nil { + logger.Error("Failed to marshal message", err, map[string]interface{}{"message": parsedMessage}) + } + msg := message.NewMessage(watermill.NewUUID(), []byte(messageText)) + err = publisher.Publish(config.Publisher.Topic, msg) + if err != nil { + logger.Error("Failed to publish message to NATS topic", err, map[string]interface{}{"topic": config.Publisher.Topic}) + return err + } + return nil +} diff --git a/internal/nats/sub.go b/internal/nats/sub.go index 2d42e3c..3c55f05 100644 --- a/internal/nats/sub.go +++ b/internal/nats/sub.go @@ -2,15 +2,16 @@ package nats import ( "caatsm/internal/config" - "caatsm/internal/handlers" "context" + "errors" "github.com/ThreeDotsLabs/watermill" "github.com/ThreeDotsLabs/watermill-nats/v2/pkg/nats" + "github.com/ThreeDotsLabs/watermill/message" nc "github.com/nats-io/nats.go" ) -func Subscribe(config *config.Config, marshaler *handlers.PlainTextMarshaler, handler *handlers.MessageHandler) { +func Subscribe(config *config.Config, marshaler *PlainTextMarshaler) { logger := watermill.NewStdLogger(false, false) options := []nc.Option{ nc.RetryOnFailedConnect(true), @@ -44,8 +45,9 @@ func Subscribe(config *config.Config, marshaler *handlers.PlainTextMarshaler, ha return } + handlers := New(config) for msg := range messages { - if err := handler.HandleMessage(msg); err == nil { + if err := handlers.HandleMessage(msg); err == nil { msg.Ack() } else { logger.Error("Failed to handle message", err, map[string]interface{}{"message": msg}) @@ -53,3 +55,17 @@ func Subscribe(config *config.Config, marshaler *handlers.PlainTextMarshaler, ha } } } + +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 +}