Added skeleton code for register endpoint (#2)
Reviewed-on: phoenix/textsender-auth#2 Co-authored-by: phoenix <kundeng00@pm.me> Co-committed-by: phoenix <kundeng00@pm.me>
This commit is contained in:
@@ -0,0 +1,113 @@
|
||||
// db/connection.go
|
||||
package db
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"log"
|
||||
"os"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"github.com/jackc/pgx/v5/pgxpool"
|
||||
)
|
||||
|
||||
var Pool *pgxpool.Pool
|
||||
|
||||
type Database struct {
|
||||
Pool *pgxpool.Pool
|
||||
}
|
||||
|
||||
func NewDatabase(connString string) (*Database, error) {
|
||||
ctx := context.Background()
|
||||
// Parse the connection string and create pool configuration
|
||||
config, err := pgxpool.ParseConfig(connString)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("unable to parse DSN: %v", err)
|
||||
}
|
||||
|
||||
// Configure connection pool settings
|
||||
config.MaxConns = 25
|
||||
config.MinConns = 5
|
||||
config.MaxConnLifetime = time.Hour
|
||||
config.MaxConnIdleTime = 30 * time.Minute
|
||||
config.HealthCheckPeriod = time.Minute
|
||||
|
||||
// Configure connection timeouts
|
||||
config.ConnConfig.ConnectTimeout = 10 * time.Second
|
||||
config.ConnConfig.RuntimeParams["timezone"] = "UTC"
|
||||
|
||||
// Create the connection pool
|
||||
pool, err := pgxpool.NewWithConfig(ctx, config)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("unable to create connection pool: %v", err)
|
||||
}
|
||||
|
||||
// Test the connection with a short timeout
|
||||
|
||||
if err := pool.Ping(ctx); err != nil {
|
||||
pool.Close()
|
||||
return nil, fmt.Errorf("unable to ping database: %v", err)
|
||||
}
|
||||
|
||||
log.Printf("Successfully connected to database")
|
||||
return &Database{Pool: pool}, nil
|
||||
}
|
||||
|
||||
func (db *Database) Close() {
|
||||
if db.Pool != nil {
|
||||
db.Pool.Close()
|
||||
log.Println("Database connection pool closed")
|
||||
}
|
||||
}
|
||||
|
||||
// HealthCheck verifies the database connection is still alive
|
||||
func (db *Database) HealthCheck(ctx context.Context) error {
|
||||
return db.Pool.Ping(ctx)
|
||||
}
|
||||
|
||||
// GetPoolStats returns connection pool statistics
|
||||
func (db *Database) GetPoolStats() string {
|
||||
stats := db.Pool.Stat()
|
||||
return fmt.Sprintf("TotalConns: %d, IdleConns: %d, AcquiredConns: %d",
|
||||
stats.TotalConns(), stats.IdleConns(), stats.AcquiredConns())
|
||||
}
|
||||
|
||||
func (db *Database) ResetDatabase(ctx context.Context) error {
|
||||
tx, err := db.Pool.Begin(ctx)
|
||||
if err != nil {
|
||||
return fmt.Errorf("Transaction unable to begin: %v", err)
|
||||
}
|
||||
defer tx.Rollback(ctx)
|
||||
|
||||
schemaContent, err := os.ReadFile("migrations/schema.sql")
|
||||
if err != nil {
|
||||
return fmt.Errorf("error reading schema file: %v", err)
|
||||
}
|
||||
|
||||
statements := strings.Split(string(schemaContent), ";")
|
||||
|
||||
for _, stmt := range statements {
|
||||
stmt = strings.TrimSpace(stmt)
|
||||
if stmt == "" {
|
||||
continue
|
||||
}
|
||||
|
||||
if strings.HasPrefix(stmt, "--") {
|
||||
continue
|
||||
}
|
||||
|
||||
_, err := tx.Exec(ctx, stmt)
|
||||
if err != nil {
|
||||
return fmt.Errorf("error executing SQL statement: %v\nStatement: %s", err, stmt)
|
||||
}
|
||||
}
|
||||
|
||||
if err := tx.Commit(ctx); err != nil {
|
||||
return fmt.Errorf("unable to commit transaction: %v", err)
|
||||
}
|
||||
|
||||
log.Println("Database schema applied successfully")
|
||||
|
||||
return nil
|
||||
}
|
||||
Reference in New Issue
Block a user