diff --git a/cmd/main/main.go b/cmd/main/main.go index d9c13da..480b8e0 100644 --- a/cmd/main/main.go +++ b/cmd/main/main.go @@ -83,6 +83,6 @@ func executeListen(c *cli.Context) error { log := utils.GetLogger() log.Info("Starting nats subscriber") // handler := handlers.New(cfg) - nats.Subscribe(cfg, &nats.PlainTextMarshaler{}) + nats.Subscribe(cfg) return nil } diff --git a/internal/nats/handler.go b/internal/nats/handler.go index 79b9dd2..4846808 100644 --- a/internal/nats/handler.go +++ b/internal/nats/handler.go @@ -53,6 +53,7 @@ func (n *MessageHandler) SaveMessage(parsed *domain.ParsedMessage, uuid string) 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}) } } diff --git a/internal/nats/sub.go b/internal/nats/sub.go index 3c55f05..cb8060f 100644 --- a/internal/nats/sub.go +++ b/internal/nats/sub.go @@ -11,8 +11,9 @@ import ( nc "github.com/nats-io/nats.go" ) -func Subscribe(config *config.Config, marshaler *PlainTextMarshaler) { +func Subscribe(config *config.Config) { logger := watermill.NewStdLogger(false, false) + marshaler := &PlainTextMarshaler{} options := []nc.Option{ nc.RetryOnFailedConnect(true), nc.Timeout(config.Timeouts.Server),