Files
cpone_midleware/internal/databaseconfig/registry.go

444 lines
14 KiB
Go

package databaseconfig
import (
"context"
"crypto/aes"
"crypto/cipher"
"crypto/rand"
"crypto/sha256"
"database/sql"
"encoding/base64"
"errors"
"fmt"
"io"
"net"
"regexp"
"sort"
"strconv"
"strings"
"sync"
"time"
mysql "github.com/go-sql-driver/mysql"
"primaya-api/cpone-middleware/internal/repository"
)
var ErrSettingNotFound = errors.New("database setting not found")
var (
validRSCode = regexp.MustCompile(`^[A-Z0-9_-]+$`)
validDatabaseName = regexp.MustCompile(`^[A-Za-z0-9_]+$`)
)
type Setting struct {
RSCode string `json:"kode_rs"`
Name string `json:"nama"`
Host string `json:"host"`
Port string `json:"port"`
Database string `json:"database"`
Username string `json:"username"`
Password string `json:"password"`
}
type PublicSetting struct {
RSCode string `json:"kode_rs"`
Name string `json:"nama"`
Host string `json:"host"`
Port string `json:"port"`
Database string `json:"database"`
Username string `json:"username"`
HasPassword bool `json:"has_password"`
Active bool `json:"active"`
}
type PoolConfig struct {
MaxOpenConns int
MaxIdleConns int
ConnMaxLifetime time.Duration
}
type Registry struct {
mu sync.RWMutex
managementDB *sql.DB
poolConfig PoolConfig
credentialKey [32]byte
settings map[string]Setting
pools map[string]*sql.DB
retiredPools []*sql.DB
}
// NewRegistry creates the management database and tables, then loads all
// active hospital database settings.
func NewRegistry(ctx context.Context, managementSetting Setting, poolConfig PoolConfig, credentialSecret string) (*Registry, error) {
if !validDatabaseName.MatchString(managementSetting.Database) {
return nil, errors.New("nama database manajemen hanya boleh berisi huruf, angka, dan underscore")
}
if strings.TrimSpace(credentialSecret) == "" {
return nil, errors.New("credential secret database wajib diisi")
}
managementDB, err := openManagementDatabase(ctx, managementSetting, poolConfig)
if err != nil {
return nil, err
}
r := &Registry{
managementDB: managementDB,
poolConfig: poolConfig,
credentialKey: sha256.Sum256([]byte(credentialSecret)),
settings: make(map[string]Setting),
pools: make(map[string]*sql.DB),
}
if err := r.migrate(ctx); err != nil {
_ = managementDB.Close()
return nil, err
}
if err := r.load(ctx); err != nil {
_ = managementDB.Close()
return nil, err
}
return r, nil
}
func NormalizeRSCode(code string) string { return strings.ToUpper(strings.TrimSpace(code)) }
func (r *Registry) Resolve(ctx context.Context, code string) (repository.MySQLLayananRepository, error) {
code = NormalizeRSCode(code)
r.mu.RLock()
setting, exists := r.settings[code]
db := r.pools[code]
r.mu.RUnlock()
if !exists {
return repository.MySQLLayananRepository{}, ErrSettingNotFound
}
if db != nil {
return repository.NewMySQLLayananRepository(db), nil
}
newDB, err := r.openAndPing(ctx, setting)
if err != nil {
return repository.MySQLLayananRepository{}, err
}
r.mu.Lock()
if existing := r.pools[code]; existing != nil {
r.mu.Unlock()
_ = newDB.Close()
return repository.NewMySQLLayananRepository(existing), nil
}
r.pools[code] = newDB
r.mu.Unlock()
return repository.NewMySQLLayananRepository(newDB), nil
}
func (r *Registry) Upsert(ctx context.Context, setting Setting) (PublicSetting, bool, error) {
setting = normalizeSetting(setting)
if err := Validate(setting); err != nil {
return PublicSetting{}, false, err
}
newDB, err := r.openAndPing(ctx, setting)
if err != nil {
return PublicSetting{}, false, fmt.Errorf("koneksi database RS gagal: %w", err)
}
encryptedPassword, err := r.encrypt(setting.Password)
if err != nil {
_ = newDB.Close()
return PublicSetting{}, false, err
}
tx, err := r.managementDB.BeginTx(ctx, nil)
if err != nil {
_ = newDB.Close()
return PublicSetting{}, false, err
}
defer tx.Rollback()
var hospitalID int64
err = tx.QueryRowContext(ctx, `SELECT id FROM hospitals WHERE code = ? FOR UPDATE`, setting.RSCode).Scan(&hospitalID)
created := errors.Is(err, sql.ErrNoRows)
if err != nil && !created {
_ = newDB.Close()
return PublicSetting{}, false, err
}
if created {
result, execErr := tx.ExecContext(ctx, `INSERT INTO hospitals (code, name, active) VALUES (?, ?, 1)`, setting.RSCode, setting.Name)
if execErr != nil {
_ = newDB.Close()
return PublicSetting{}, false, execErr
}
hospitalID, err = result.LastInsertId()
if err != nil {
_ = newDB.Close()
return PublicSetting{}, false, err
}
} else if _, err = tx.ExecContext(ctx, `UPDATE hospitals SET name = ?, active = 1 WHERE id = ?`, setting.Name, hospitalID); err != nil {
_ = newDB.Close()
return PublicSetting{}, false, err
}
_, err = tx.ExecContext(ctx, `
INSERT INTO hospital_databases (hospital_id, host, port, database_name, username, encrypted_password)
VALUES (?, ?, ?, ?, ?, ?)
ON DUPLICATE KEY UPDATE host = VALUES(host), port = VALUES(port), database_name = VALUES(database_name),
username = VALUES(username), encrypted_password = VALUES(encrypted_password)`,
hospitalID, setting.Host, setting.Port, setting.Database, setting.Username, encryptedPassword)
if err != nil {
_ = newDB.Close()
return PublicSetting{}, false, err
}
if err = tx.Commit(); err != nil {
_ = newDB.Close()
return PublicSetting{}, false, err
}
r.mu.Lock()
oldPool := r.pools[setting.RSCode]
r.settings[setting.RSCode] = setting
r.pools[setting.RSCode] = newDB
if oldPool != nil {
r.retiredPools = append(r.retiredPools, oldPool)
}
r.mu.Unlock()
return publicSetting(setting), created, nil
}
func (r *Registry) List() []PublicSetting {
r.mu.RLock()
defer r.mu.RUnlock()
result := make([]PublicSetting, 0, len(r.settings))
for _, setting := range r.settings {
result = append(result, publicSetting(setting))
}
sort.Slice(result, func(i, j int) bool { return result[i].RSCode < result[j].RSCode })
return result
}
func (r *Registry) Get(code string) (PublicSetting, bool) {
code = NormalizeRSCode(code)
r.mu.RLock()
defer r.mu.RUnlock()
setting, exists := r.settings[code]
return publicSetting(setting), exists
}
func (r *Registry) Close() error {
r.mu.Lock()
defer r.mu.Unlock()
var firstErr error
for code, db := range r.pools {
if err := db.Close(); err != nil && firstErr == nil {
firstErr = err
}
delete(r.pools, code)
}
for _, db := range r.retiredPools {
if err := db.Close(); err != nil && firstErr == nil {
firstErr = err
}
}
if err := r.managementDB.Close(); err != nil && firstErr == nil {
firstErr = err
}
r.retiredPools = nil
return firstErr
}
func openManagementDatabase(ctx context.Context, setting Setting, pool PoolConfig) (*sql.DB, error) {
serverSetting := setting
serverSetting.Database = ""
bootstrap, err := sql.Open("mysql", mysqlDSN(serverSetting))
if err != nil {
return nil, err
}
if err := bootstrap.PingContext(ctx); err != nil {
_ = bootstrap.Close()
return nil, fmt.Errorf("koneksi server database manajemen gagal: %w", err)
}
_, err = bootstrap.ExecContext(ctx, "CREATE DATABASE IF NOT EXISTS `"+setting.Database+"` CHARACTER SET utf8mb4 COLLATE utf8mb4_unicode_ci")
_ = bootstrap.Close()
if err != nil {
return nil, fmt.Errorf("buat database manajemen %s: %w", setting.Database, err)
}
db, err := sql.Open("mysql", mysqlDSN(setting))
if err != nil {
return nil, err
}
applyPoolConfig(db, pool)
if err := db.PingContext(ctx); err != nil {
_ = db.Close()
return nil, fmt.Errorf("koneksi database manajemen gagal: %w", err)
}
return db, nil
}
func (r *Registry) migrate(ctx context.Context) error {
statements := []string{
`CREATE TABLE IF NOT EXISTS hospitals (
id BIGINT UNSIGNED NOT NULL AUTO_INCREMENT, code VARCHAR(50) NOT NULL, name VARCHAR(150) NOT NULL,
active TINYINT(1) NOT NULL DEFAULT 1,
created_at TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP,
updated_at TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
PRIMARY KEY (id), UNIQUE KEY uq_hospitals_code (code), KEY idx_hospitals_active (active)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_unicode_ci`,
`CREATE TABLE IF NOT EXISTS hospital_databases (
id BIGINT UNSIGNED NOT NULL AUTO_INCREMENT, hospital_id BIGINT UNSIGNED NOT NULL,
host VARCHAR(255) NOT NULL, port SMALLINT UNSIGNED NOT NULL DEFAULT 3306,
database_name VARCHAR(100) NOT NULL, username VARCHAR(100) NOT NULL, encrypted_password TEXT NOT NULL,
created_at TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP,
updated_at TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
PRIMARY KEY (id), UNIQUE KEY uq_hospital_databases_hospital (hospital_id),
CONSTRAINT fk_hospital_databases_hospital FOREIGN KEY (hospital_id) REFERENCES hospitals(id) ON DELETE CASCADE
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_unicode_ci`,
}
for _, statement := range statements {
if _, err := r.managementDB.ExecContext(ctx, statement); err != nil {
return fmt.Errorf("migrasi database manajemen: %w", err)
}
}
var defaultColumnExists int
if err := r.managementDB.QueryRowContext(ctx, `SELECT COUNT(*) FROM information_schema.columns
WHERE table_schema = DATABASE() AND table_name = 'hospitals' AND column_name = 'is_default'`).Scan(&defaultColumnExists); err != nil {
return fmt.Errorf("periksa schema database manajemen: %w", err)
}
if defaultColumnExists > 0 {
if _, err := r.managementDB.ExecContext(ctx, `ALTER TABLE hospitals DROP COLUMN is_default`); err != nil {
return fmt.Errorf("hapus kolom default rumah sakit: %w", err)
}
}
return nil
}
func (r *Registry) load(ctx context.Context) error {
rows, err := r.managementDB.QueryContext(ctx, `SELECT h.code, h.name, d.host, d.port, d.database_name, d.username, d.encrypted_password
FROM hospitals h INNER JOIN hospital_databases d ON d.hospital_id = h.id WHERE h.active = 1 ORDER BY h.code`)
if err != nil {
return err
}
defer rows.Close()
settings := make(map[string]Setting)
for rows.Next() {
var setting Setting
var encryptedPassword string
if err := rows.Scan(&setting.RSCode, &setting.Name, &setting.Host, &setting.Port, &setting.Database, &setting.Username, &encryptedPassword); err != nil {
return err
}
setting.Password, err = r.decrypt(encryptedPassword)
if err != nil {
return fmt.Errorf("decrypt password %s: %w", setting.RSCode, err)
}
settings[setting.RSCode] = setting
}
if err := rows.Err(); err != nil {
return err
}
r.settings = settings
return nil
}
func (r *Registry) encrypt(plainText string) (string, error) {
block, err := aes.NewCipher(r.credentialKey[:])
if err != nil {
return "", err
}
gcm, err := cipher.NewGCM(block)
if err != nil {
return "", err
}
nonce := make([]byte, gcm.NonceSize())
if _, err := io.ReadFull(rand.Reader, nonce); err != nil {
return "", err
}
return base64.StdEncoding.EncodeToString(gcm.Seal(nonce, nonce, []byte(plainText), nil)), nil
}
func (r *Registry) decrypt(encoded string) (string, error) {
sealed, err := base64.StdEncoding.DecodeString(encoded)
if err != nil {
return "", err
}
block, err := aes.NewCipher(r.credentialKey[:])
if err != nil {
return "", err
}
gcm, err := cipher.NewGCM(block)
if err != nil {
return "", err
}
if len(sealed) < gcm.NonceSize() {
return "", errors.New("encrypted password tidak valid")
}
plainText, err := gcm.Open(nil, sealed[:gcm.NonceSize()], sealed[gcm.NonceSize():], nil)
if err != nil {
return "", err
}
return string(plainText), nil
}
func (r *Registry) openAndPing(ctx context.Context, setting Setting) (*sql.DB, error) {
db, err := sql.Open("mysql", mysqlDSN(setting))
if err != nil {
return nil, err
}
applyPoolConfig(db, r.poolConfig)
if err := db.PingContext(ctx); err != nil {
_ = db.Close()
return nil, err
}
return db, nil
}
func applyPoolConfig(db *sql.DB, pool PoolConfig) {
db.SetMaxOpenConns(pool.MaxOpenConns)
db.SetMaxIdleConns(pool.MaxIdleConns)
db.SetConnMaxLifetime(pool.ConnMaxLifetime)
}
func Validate(setting Setting) error {
setting = normalizeSetting(setting)
if setting.RSCode == "" {
return errors.New("kode_rs wajib diisi")
}
if len(setting.RSCode) > 50 {
return errors.New("kode_rs maksimal 50 karakter")
}
if !validRSCode.MatchString(setting.RSCode) {
return errors.New("kode_rs hanya boleh berisi huruf, angka, tanda hubung, dan underscore")
}
if setting.Name == "" {
return errors.New("nama rumah sakit wajib diisi")
}
if setting.Host == "" || setting.Port == "" || setting.Database == "" || setting.Username == "" {
return errors.New("host, port, database, dan username wajib diisi")
}
port, err := strconv.Atoi(setting.Port)
if err != nil || port < 1 || port > 65535 {
return errors.New("port harus berupa angka antara 1 sampai 65535")
}
return nil
}
func normalizeSetting(setting Setting) Setting {
setting.RSCode = NormalizeRSCode(setting.RSCode)
setting.Name = strings.TrimSpace(setting.Name)
setting.Host = strings.TrimSpace(setting.Host)
setting.Port = strings.TrimSpace(setting.Port)
setting.Database = strings.TrimSpace(setting.Database)
setting.Username = strings.TrimSpace(setting.Username)
return setting
}
func publicSetting(setting Setting) PublicSetting {
return PublicSetting{RSCode: setting.RSCode, Name: setting.Name, Host: setting.Host, Port: setting.Port,
Database: setting.Database, Username: setting.Username, HasPassword: setting.Password != "",
Active: true}
}
func mysqlDSN(setting Setting) string {
cfg := mysql.NewConfig()
cfg.User = setting.Username
cfg.Passwd = setting.Password
cfg.Net = "tcp"
cfg.Addr = net.JoinHostPort(setting.Host, setting.Port)
cfg.DBName = setting.Database
cfg.ParseTime = true
cfg.Loc = time.Local
cfg.Params = map[string]string{"charset": "utf8mb4"}
return cfg.FormatDSN()
}