diff --git a/internal/iface/interface.go b/internal/iface/interface.go index 1ea174b..9bd48f3 100644 --- a/internal/iface/interface.go +++ b/internal/iface/interface.go @@ -6,7 +6,7 @@ import ( ) type MessageHandler interface { - HandleMessage(msg []byte, uuid []byte) error + HandleMessage(msg []byte, id string) error } type MessagePublisher interface { @@ -18,5 +18,5 @@ type MessageSubscriber interface { } type MessageRepository interface { - CreateNew(message *domain.ParsedMessage, uuid []byte) error + CreateNew(message *domain.ParsedMessage) error } diff --git a/internal/nats/handler.go b/internal/nats/handler.go index f977904..781c6a9 100644 --- a/internal/nats/handler.go +++ b/internal/nats/handler.go @@ -25,7 +25,7 @@ func NewHandler(config *config.Config, publisher iface.MessagePublisher, reposit } } -func (handler *MessageHandler) HandleMessage(msg []byte, uuid []byte) error { +func (handler *MessageHandler) HandleMessage(msg []byte, id string) error { handler.mu.Lock() defer handler.mu.Unlock() log := utils.GetSugaredLogger() @@ -36,28 +36,14 @@ func (handler *MessageHandler) HandleMessage(msg []byte, uuid []byte) error { payload := string(msg) var parsed *domain.ParsedMessage if parsed = parsers.Parse(payload); !parsed.Parsed { - log.Infof("not parsed: [%s] : {%s} \n", uuid, payload) + log.Infof("not parsed: [%s] : {%s} \n", id, payload) } else { - log.Infof("parsed [%s]: %v\n", uuid, parsed.ToString()) + parsed.Uuid = id + log.Infof("parsed [%s]: %v\n", id, parsed.ToString()) } - handler.repository.CreateNew(parsed, uuid) + handler.repository.CreateNew(parsed) + handler.publisher.Publish(parsed) + return nil } - -// func (n *MessageHandler) SaveMessage(parsed *domain.ParsedMessage, uuid string) { -// logger := watermill.NewStdLogger(false, false) -// if parsed != nil { -// parsed.Uuid = uuid -// 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}) -// } -// } - -// 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}) -// } -// } diff --git a/internal/nats/pub.go b/internal/nats/pub.go index 9a591d5..574042c 100644 --- a/internal/nats/pub.go +++ b/internal/nats/pub.go @@ -50,5 +50,6 @@ func (n *NatsPublisher) Publish(parsedMessage interface{}) error { logger.Errorf("Failed to publish message: %v", err) return err } + logger.Infof("Message published: %s", msg.UUID) return nil } diff --git a/internal/nats/sub.go b/internal/nats/sub.go index 4dbe7f9..56498c2 100644 --- a/internal/nats/sub.go +++ b/internal/nats/sub.go @@ -55,10 +55,11 @@ func (n *NatsSubscriber) Subscribe(config *config.Config, handlers iface.Message } for msg := range messages { - if err := handlers.HandleMessage(msg.Payload, []byte(msg.UUID)); err == nil { + if err := handlers.HandleMessage(msg.Payload, msg.UUID); err == nil { + logger.Infof("Message handled: %s", msg.UUID) msg.Ack() } else { - logger.Error("Failed to handle message", err, map[string]interface{}{"message": msg}) + logger.Errorf("Failed to handle message [%s]: %v", msg.UUID, err) msg.Nack() } } diff --git a/internal/parsers/aviation_parser.go b/internal/parsers/aviation_parser.go index 7afb50b..fb5a80a 100644 --- a/internal/parsers/aviation_parser.go +++ b/internal/parsers/aviation_parser.go @@ -8,6 +8,8 @@ import ( "strings" "sync" "time" + + "github.com/google/uuid" ) const ( @@ -211,6 +213,7 @@ func Parse(rawText string) *domain.ParsedMessage { } message.Parsed = true message.BodyData = bodyData + message.Uuid = uuid.New().String() return &message } diff --git a/internal/repository/hasura.go b/internal/repository/hasura.go index abe2d2b..8ae82f6 100644 --- a/internal/repository/hasura.go +++ b/internal/repository/hasura.go @@ -10,7 +10,6 @@ import ( "caatsm/pkg/utils" "github.com/Khan/genqlient/graphql" - "github.com/google/uuid" "golang.org/x/oauth2" ) @@ -34,10 +33,12 @@ func NewHasura(config *config.Config) *HasuraRepository { } // InsertParsedMessage inserts a new ParsedMessage into the Hasura GraphQL API -func (hr *HasuraRepository) CreateNew(pm *domain.ParsedMessage, msg_uuid []byte) error { +func (hr *HasuraRepository) CreateNew(pm *domain.ParsedMessage) error { log := utils.GetSugaredLogger() bodyString, _ := json.Marshal(pm.BodyData) secondAddress, _ := json.Marshal(pm.SecondaryAddresses) + var err error + msgUuid := utils.GetUuid(pm.Uuid) variables := Aviation_telegrams_insert_input{ Message_id: pm.MessageID, Priority_indicator: pm.PriorityIndicator, @@ -48,7 +49,7 @@ func (hr *HasuraRepository) CreateNew(pm *domain.ParsedMessage, msg_uuid []byte) Category: pm.Category, Date_time: pm.DateTime, Dispatched_at: pm.DispatchedAt, - Uuid: uuid.Must(uuid.FromBytes(msg_uuid)), + Uuid: msgUuid, Received_at: pm.ReceivedAt, Originator: pm.Originator, Originator_date_time: pm.OriginatorDateTime, diff --git a/pkg/utils/util.go b/pkg/utils/util.go new file mode 100644 index 0000000..c54846f --- /dev/null +++ b/pkg/utils/util.go @@ -0,0 +1,13 @@ +package utils + +import "github.com/google/uuid" + +func GetUuid(uuidString string) uuid.UUID { + log := GetSugaredLogger() + result, err := uuid.Parse(uuidString) + if err != nil { + log.Warnf("invalid uuid: %s", uuidString) + return uuid.New() + } + return result +}