refactor: Update main.go to process Civial Aviation Authority Telegram messages and add CLI flag for Nats server address and topic
This commit is contained in:
+45
-19
@@ -12,6 +12,10 @@ import (
|
|||||||
"github.com/urfave/cli/v2"
|
"github.com/urfave/cli/v2"
|
||||||
)
|
)
|
||||||
|
|
||||||
|
var (
|
||||||
|
cfg *config.Config
|
||||||
|
)
|
||||||
|
|
||||||
func main() {
|
func main() {
|
||||||
app := setupApp()
|
app := setupApp()
|
||||||
if err := app.Run(os.Args); err != nil {
|
if err := app.Run(os.Args); err != nil {
|
||||||
@@ -21,26 +25,10 @@ func main() {
|
|||||||
|
|
||||||
func setupApp() *cli.App {
|
func setupApp() *cli.App {
|
||||||
app := &cli.App{
|
app := &cli.App{
|
||||||
Name: "serial-read",
|
Name: "telegram message process",
|
||||||
Usage: "A serial port reading CLI application",
|
Usage: "A Civial Aviation Authority Telegram message processor",
|
||||||
Before: func(c *cli.Context) error {
|
Before: func(c *cli.Context) error {
|
||||||
// parameter = config.GetParameter()
|
|
||||||
return nil
|
|
||||||
},
|
|
||||||
Commands: []*cli.Command{
|
|
||||||
{
|
|
||||||
Name: "listen",
|
|
||||||
Usage: "Listen to nats messages",
|
|
||||||
Action: executeListen,
|
|
||||||
},
|
|
||||||
},
|
|
||||||
}
|
|
||||||
return app
|
|
||||||
}
|
|
||||||
|
|
||||||
func executeListen(c *cli.Context) error {
|
|
||||||
cfg, err := config.LoadConfig()
|
cfg, err := config.LoadConfig()
|
||||||
log := utils.GetLogger()
|
|
||||||
if err != nil {
|
if err != nil {
|
||||||
fmt.Printf("Error loading configuration: %v\n", err)
|
fmt.Printf("Error loading configuration: %v\n", err)
|
||||||
return err
|
return err
|
||||||
@@ -50,9 +38,47 @@ func executeListen(c *cli.Context) error {
|
|||||||
fmt.Printf("Invalid configuration: %v\n", err)
|
fmt.Printf("Invalid configuration: %v\n", err)
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
|
overrideConfig(c)
|
||||||
|
return nil
|
||||||
|
},
|
||||||
|
Commands: []*cli.Command{
|
||||||
|
{
|
||||||
|
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,
|
||||||
|
},
|
||||||
|
},
|
||||||
|
}
|
||||||
|
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()
|
||||||
|
log := utils.GetLogger()
|
||||||
fmt.Println("Loaded configuration successfully")
|
fmt.Println("Loaded configuration successfully")
|
||||||
|
|
||||||
log.Info("Starting nats subscriber")
|
log.Info("Starting nats subscriber")
|
||||||
handler := handlers.NewNatsHandler(cfg)
|
handler := handlers.NewNatsHandler(cfg)
|
||||||
|
|
||||||
|
|||||||
@@ -15,17 +15,17 @@ import (
|
|||||||
nc "github.com/nats-io/nats.go"
|
nc "github.com/nats-io/nats.go"
|
||||||
)
|
)
|
||||||
|
|
||||||
type NatsHandler struct {
|
type MessageHandler struct {
|
||||||
mu sync.Mutex
|
mu sync.Mutex
|
||||||
config *config.Config
|
config *config.Config
|
||||||
hasuraRepo *repository.HasuraRepository
|
hasuraRepo *repository.HasuraRepository
|
||||||
}
|
}
|
||||||
|
|
||||||
func NewNatsHandler(config *config.Config) *NatsHandler {
|
func NewNatsHandler(config *config.Config) *MessageHandler {
|
||||||
return &NatsHandler{config: config, hasuraRepo: repository.NewHasuraRepo(config.Hasura.Endpoint, config.Hasura.Secret)}
|
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()
|
n.mu.Lock()
|
||||||
defer n.mu.Unlock()
|
defer n.mu.Unlock()
|
||||||
log := utils.GetSugaredLogger()
|
log := utils.GetSugaredLogger()
|
||||||
@@ -45,7 +45,7 @@ func (n *NatsHandler) HandleMessage(msg *message.Message) error {
|
|||||||
return nil
|
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)
|
logger := watermill.NewStdLogger(false, false)
|
||||||
if parsed != nil {
|
if parsed != nil {
|
||||||
parsed.Uuid = uuid
|
parsed.Uuid = uuid
|
||||||
@@ -10,7 +10,7 @@ import (
|
|||||||
nc "github.com/nats-io/nats.go"
|
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{}
|
// marshaler := &PlainTextMarshaler{}
|
||||||
logger := watermill.NewStdLogger(false, false)
|
logger := watermill.NewStdLogger(false, false)
|
||||||
options := []nc.Option{
|
options := []nc.Option{
|
||||||
|
|||||||
Reference in New Issue
Block a user