diff --git a/cmd/main/main.go b/cmd/main/main.go index 5766644..586072a 100644 --- a/cmd/main/main.go +++ b/cmd/main/main.go @@ -12,6 +12,10 @@ import ( "github.com/urfave/cli/v2" ) +var ( + cfg *config.Config +) + func main() { app := setupApp() if err := app.Run(os.Args); err != nil { @@ -21,16 +25,40 @@ func main() { func setupApp() *cli.App { app := &cli.App{ - Name: "serial-read", - Usage: "A serial port reading CLI application", + Name: "telegram message process", + Usage: "A Civial Aviation Authority Telegram message processor", Before: func(c *cli.Context) error { - // parameter = config.GetParameter() + cfg, err := config.LoadConfig() + if err != nil { + fmt.Printf("Error loading configuration: %v\n", err) + return err + } + + if err := config.ValidateConfig(cfg); err != nil { + fmt.Printf("Invalid configuration: %v\n", err) + return err + } + overrideConfig(c) return nil }, Commands: []*cli.Command{ { - Name: "listen", - Usage: "Listen to nats messages", + Name: "listen", + Usage: "Listen to nats messages", + Flags: []cli.Flag{ + &cli.StringFlag{ + Name: "nats", + Usage: "Nats server address", + Value: "nats://localhost:4222", + EnvVars: []string{"NATS_SERVER"}, + }, + &cli.StringFlag{ + Name: "topic", + Usage: "Nats topic to listen to", + Value: "Telegram.Serial", + EnvVars: []string{"NATS_SUBJECT"}, + }, + }, Action: executeListen, }, }, @@ -38,21 +66,19 @@ func setupApp() *cli.App { return app } +func overrideConfig(c *cli.Context) { + if c.IsSet("nats") { + cfg.Nats.URL = c.String("nats") + } + if c.IsSet("topic") { + cfg.Subscription.Topic = c.String("topic") + } +} + func executeListen(c *cli.Context) error { - cfg, err := config.LoadConfig() + // cfg, err := config.LoadConfig() log := utils.GetLogger() - if err != nil { - fmt.Printf("Error loading configuration: %v\n", err) - return err - } - - if err := config.ValidateConfig(cfg); err != nil { - fmt.Printf("Invalid configuration: %v\n", err) - return err - } - fmt.Println("Loaded configuration successfully") - log.Info("Starting nats subscriber") handler := handlers.NewNatsHandler(cfg) diff --git a/internal/handlers/nats_handler.go b/internal/handlers/message_handler.go similarity index 80% rename from internal/handlers/nats_handler.go rename to internal/handlers/message_handler.go index da002a1..f9c421f 100644 --- a/internal/handlers/nats_handler.go +++ b/internal/handlers/message_handler.go @@ -15,17 +15,17 @@ import ( nc "github.com/nats-io/nats.go" ) -type NatsHandler struct { +type MessageHandler struct { mu sync.Mutex config *config.Config hasuraRepo *repository.HasuraRepository } -func NewNatsHandler(config *config.Config) *NatsHandler { - return &NatsHandler{config: config, hasuraRepo: repository.NewHasuraRepo(config.Hasura.Endpoint, config.Hasura.Secret)} +func NewNatsHandler(config *config.Config) *MessageHandler { + return &MessageHandler{config: config, hasuraRepo: repository.NewHasuraRepo(config.Hasura.Endpoint, config.Hasura.Secret)} } -func (n *NatsHandler) HandleMessage(msg *message.Message) error { +func (n *MessageHandler) HandleMessage(msg *message.Message) error { n.mu.Lock() defer n.mu.Unlock() log := utils.GetSugaredLogger() @@ -45,7 +45,7 @@ func (n *NatsHandler) HandleMessage(msg *message.Message) error { return nil } -func (n *NatsHandler) SaveMessage(parsed *domain.ParsedMessage, uuid string) { +func (n *MessageHandler) SaveMessage(parsed *domain.ParsedMessage, uuid string) { logger := watermill.NewStdLogger(false, false) if parsed != nil { parsed.Uuid = uuid diff --git a/internal/nats/sub.go b/internal/nats/sub.go index 443abde..41cab76 100644 --- a/internal/nats/sub.go +++ b/internal/nats/sub.go @@ -10,7 +10,7 @@ import ( nc "github.com/nats-io/nats.go" ) -func Subscribe(config *config.Config, marshaler *handlers.PlainTextMarshaler, handler *handlers.NatsHandler) { +func Subscribe(config *config.Config, marshaler *handlers.PlainTextMarshaler, handler *handlers.MessageHandler) { // marshaler := &PlainTextMarshaler{} logger := watermill.NewStdLogger(false, false) options := []nc.Option{