feat: add dynamic hospital database registry
This commit is contained in:
484
internal/databaseconfig/registry.go
Normal file
484
internal/databaseconfig/registry.go
Normal file
@@ -0,0 +1,484 @@
|
||||
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"`
|
||||
IsDefault bool `json:"is_default"`
|
||||
Active bool `json:"active"`
|
||||
}
|
||||
|
||||
type PoolConfig struct {
|
||||
MaxOpenConns int
|
||||
MaxIdleConns int
|
||||
ConnMaxLifetime time.Duration
|
||||
}
|
||||
|
||||
type Registry struct {
|
||||
mu sync.RWMutex
|
||||
managementDB *sql.DB
|
||||
defaultCode string
|
||||
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 seeds the
|
||||
// existing HIS connection as the default hospital.
|
||||
func NewRegistry(ctx context.Context, managementSetting Setting, defaultCode string, defaultSetting 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,
|
||||
defaultCode: NormalizeRSCode(defaultCode),
|
||||
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
|
||||
}
|
||||
defaultSetting.RSCode = r.defaultCode
|
||||
if strings.TrimSpace(defaultSetting.Name) == "" {
|
||||
defaultSetting.Name = "RS Dev Awalbros"
|
||||
}
|
||||
if err := r.seedDefault(ctx, defaultSetting); 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) DefaultCode() string { return r.defaultCode }
|
||||
|
||||
func (r *Registry) Resolve(ctx context.Context, code string) (repository.MySQLLayananRepository, error) {
|
||||
code = NormalizeRSCode(code)
|
||||
if code == "" {
|
||||
code = r.defaultCode
|
||||
}
|
||||
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, is_default, active) VALUES (?, ?, 0, 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, r.defaultCode), 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, r.defaultCode))
|
||||
}
|
||||
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, r.defaultCode), 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,
|
||||
is_default TINYINT(1) NOT NULL DEFAULT 0, 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)
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (r *Registry) seedDefault(ctx context.Context, setting Setting) error {
|
||||
setting = normalizeSetting(setting)
|
||||
if err := Validate(setting); err != nil {
|
||||
return fmt.Errorf("setting database default tidak valid: %w", err)
|
||||
}
|
||||
encryptedPassword, err := r.encrypt(setting.Password)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
tx, err := r.managementDB.BeginTx(ctx, nil)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
defer tx.Rollback()
|
||||
if _, err := tx.ExecContext(ctx, `UPDATE hospitals SET is_default = 0 WHERE code <> ?`, setting.RSCode); err != nil {
|
||||
return err
|
||||
}
|
||||
_, err = tx.ExecContext(ctx, `INSERT INTO hospitals (code, name, is_default, active) VALUES (?, ?, 1, 1)
|
||||
ON DUPLICATE KEY UPDATE is_default = 1, active = 1`, setting.RSCode, setting.Name)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
var hospitalID int64
|
||||
if err := tx.QueryRowContext(ctx, `SELECT id FROM hospitals WHERE code = ?`, setting.RSCode).Scan(&hospitalID); err != nil {
|
||||
return err
|
||||
}
|
||||
_, err = tx.ExecContext(ctx, `INSERT INTO hospital_databases (hospital_id, host, port, database_name, username, encrypted_password)
|
||||
VALUES (?, ?, ?, ?, ?, ?) ON DUPLICATE KEY UPDATE hospital_id = hospital_id`, hospitalID, setting.Host,
|
||||
setting.Port, setting.Database, setting.Username, encryptedPassword)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
return tx.Commit()
|
||||
}
|
||||
|
||||
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, defaultCode string) PublicSetting {
|
||||
return PublicSetting{RSCode: setting.RSCode, Name: setting.Name, Host: setting.Host, Port: setting.Port,
|
||||
Database: setting.Database, Username: setting.Username, HasPassword: setting.Password != "",
|
||||
IsDefault: setting.RSCode == defaultCode, 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()
|
||||
}
|
||||
Reference in New Issue
Block a user