package db import ( "context" "database/sql" "fmt" "os" "time" "doormile/config" "doormile/utils" nats "github.com/nats-io/nats.go" "github.com/redis/go-redis/v9" "gorm.io/driver/postgres" "gorm.io/gorm" "gorm.io/gorm/logger" ) var ( DB *gorm.DB Rdb *redis.Client Ctx = context.Background() Nc *nats.Conn Js nats.JetStreamContext ) func Connect(cfg *config.Config) { dsn := fmt.Sprintf( "host=%s user=%s password=%s dbname=%s port=%s sslmode=disable TimeZone=Asia/Kolkata", cfg.DBHost, cfg.DBUser, cfg.DBPassword, cfg.DBName, cfg.DBPort, ) var err error maxRetries := 5 backoff := 2 * time.Second for i := 1; i <= maxRetries; i++ { utils.Info("Connecting to database", "attempt", i, "host", cfg.DBHost, "port", cfg.DBPort) DB, err = gorm.Open(postgres.Open(dsn), &gorm.Config{ Logger: logger.Default.LogMode(logger.Error), }) if err == nil { sqlDB, dbErr := DB.DB() if dbErr == nil { ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) defer cancel() if pingErr := sqlDB.PingContext(ctx); pingErr == nil { setupDB(DB) utils.Info("✅ Database connected successfully") startDBHealthCheck(sqlDB) return } else { err = pingErr } } else { err = dbErr } } utils.Error("❌ DB connection failed", "attempt", i, "error", err) time.Sleep(backoff) backoff *= 2 } utils.Logger.Fatal("☢️ CRITICAL: DB connection failed after retries") } func setupDB(database *gorm.DB) { sqlDB, err := database.DB() if err != nil { utils.Error("Failed to get sql.DB", "error", err) return } sqlDB.SetMaxOpenConns(30) sqlDB.SetMaxIdleConns(5) sqlDB.SetConnMaxLifetime(5 * time.Minute) sqlDB.SetConnMaxIdleTime(2 * time.Minute) } func startDBHealthCheck(sqlDB *sql.DB) { go func() { for { ctx, cancel := context.WithTimeout(context.Background(), 3*time.Second) err := sqlDB.PingContext(ctx) cancel() if err != nil { utils.Error("❌ DB connection lost", "error", err) } time.Sleep(30 * time.Second) } }() } func CloseDB() { if DB == nil { return } sqlDB, err := DB.DB() if err != nil { utils.Error("Error retrieving sql.DB for close", "error", err) return } utils.Info("Closing DB connection") sqlDB.Close() } func InitRedis(cfg *config.Config) { addr := fmt.Sprintf("%s:%s", cfg.RedisHost, cfg.RedisPort) Rdb = redis.NewClient(&redis.Options{ Addr: addr, Username: cfg.RedisUser, Password: cfg.RedisPassword, DB: 0, DialTimeout: 5 * time.Second, ReadTimeout: 3 * time.Second, WriteTimeout: 3 * time.Second, PoolSize: 10, MinIdleConns: 2, }) ctx, cancel := context.WithTimeout(context.Background(), 3*time.Second) defer cancel() _, err := Rdb.Ping(ctx).Result() if err != nil { utils.Error("⚠️ Redis connection failed, continuing in fallback mode", "addr", addr, "error", err) } else { utils.Info("✅ Redis connected successfully", "addr", addr) } } func InitNATS(cfg *config.Config) { var err error Nc, err = nats.Connect( cfg.NatsURL, nats.UserInfo(cfg.NatsUser, cfg.NatsPassword), nats.ReconnectWait(2*time.Second), nats.MaxReconnects(-1), ) if err != nil { utils.Error("⚠️ NATS connection failed, continuing without NATS", "url", cfg.NatsURL, "error", err) return } Js, err = Nc.JetStream() if err != nil { utils.Error("⚠️ NATS JetStream init failed", "error", err) Nc.Close() Nc = nil return } utils.Info("✅ NATS JetStream connected successfully", "url", cfg.NatsURL) } func getEnv(key, fallback string) string { if val := os.Getenv(key); val != "" { return val } return fallback }