107 lines
2.1 KiB
Go
107 lines
2.1 KiB
Go
package main
|
|
|
|
import (
|
|
"context"
|
|
"log"
|
|
"os"
|
|
"os/signal"
|
|
"smpp-transmitter/config"
|
|
"smpp-transmitter/pkg/data"
|
|
fl "smpp-transmitter/pkg/logger"
|
|
"smpp-transmitter/pkg/mq"
|
|
"smpp-transmitter/pkg/mq/rabbitmq"
|
|
"syscall"
|
|
|
|
"github.com/joho/godotenv"
|
|
"gorm.io/gorm"
|
|
)
|
|
|
|
func main() {
|
|
if err := run(); err != nil {
|
|
log.Fatalf("Application error: %v", err)
|
|
}
|
|
}
|
|
|
|
func run() error {
|
|
// Load configuration
|
|
cfg, err := loadConfig()
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
// Setup logger
|
|
logger, err := fl.SetupLogger(cfg.AppEnv)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
// Initialize database
|
|
db, err := data.InitDB(cfg.MysqlDSN)
|
|
if err != nil {
|
|
logger.Error("Failed to connect to MySQL:", err)
|
|
return err
|
|
}
|
|
|
|
// Initialize RabbitMQ consumer
|
|
consumer, err := initRabbitMQConsumer(cfg)
|
|
if err != nil {
|
|
logger.Error("Failed to initialize RabbitMQ consumer:", err)
|
|
return err
|
|
}
|
|
|
|
// Start consuming messages
|
|
go startConsuming(consumer, db, logger)
|
|
|
|
// Handle graceful shutdown
|
|
waitForShutdown(logger)
|
|
|
|
return nil
|
|
}
|
|
|
|
func loadConfig() (*config.Config, error) {
|
|
if err := godotenv.Load(); err != nil {
|
|
return nil, err
|
|
}
|
|
|
|
return config.NewConfig()
|
|
}
|
|
|
|
func initRabbitMQConsumer(cfg *config.Config) (mq.Consumer, error) {
|
|
return rabbitmq.NewRabbitMQConsumer(
|
|
cfg.RabbitMQURL,
|
|
cfg.ExchangeName,
|
|
cfg.QueueName,
|
|
cfg.RoutingKey,
|
|
1,
|
|
), nil
|
|
}
|
|
|
|
func startConsuming(consumer mq.Consumer, db *gorm.DB, logger fl.LogService) {
|
|
consumer.Consume(func(ctx context.Context, body []byte) {
|
|
handleMessage(ctx, body, db, logger)
|
|
})
|
|
}
|
|
|
|
func handleMessage(ctx context.Context, body []byte, db *gorm.DB, logger fl.LogService) {
|
|
message := &data.Message{}
|
|
|
|
if err := message.Convert(body); err != nil {
|
|
logger.Error("Message cannot be converted", err)
|
|
return
|
|
}
|
|
|
|
logger.Info("Received message", message.Src, message.Dst, message.Msg)
|
|
|
|
if err := message.Insert(ctx, db); err != nil {
|
|
logger.Error("Message cannot be Inserted", err)
|
|
return
|
|
}
|
|
}
|
|
|
|
func waitForShutdown(logger fl.LogService) {
|
|
sigs := make(chan os.Signal, 1)
|
|
signal.Notify(sigs, syscall.SIGINT, syscall.SIGTERM)
|
|
<-sigs
|
|
logger.Info("Shutting down...")
|
|
}
|