kra-new/app/system/internal/data/data.go

299 lines
9.2 KiB
Go

package data
import (
"context"
"fmt"
"log/slog"
"sync"
"sync/atomic"
"time"
"github.com/google/wire"
"github.com/redis/go-redis/v9"
"gorm.io/gorm"
"kra/app/system/internal/conf"
datapayment "kra/app/system/internal/data/payment"
datasystem "kra/app/system/internal/data/repository"
"kra/app/system/internal/integration/storage"
"kra/pkg/module"
)
var ProviderSet = wire.NewSet(
NewData,
wire.Bind(new(datasystem.Provider), new(*Data)),
wire.Bind(new(datasystem.DatabaseProvider), new(*Data)),
wire.Bind(new(datapayment.Provider), new(*Data)),
datasystem.NewRuntimeSettings,
datasystem.NewTokenIssuer,
datasystem.NewUserRepo, datasystem.NewAuthorityAccessRepo, datasystem.NewAPIRepo, datasystem.NewPermissionRepo,
datasystem.NewMenuRepo, datasystem.NewDepartmentRepo, datasystem.NewPositionRepo, datasystem.NewDictionaryRepo, datasystem.NewParameterRepo, datasystem.NewAPITokenRepo,
datasystem.NewSecurityRepo,
datasystem.NewVersionRepo, datasystem.NewExportRepo, datasystem.NewAuditRepo, datasystem.NewAuditRecorderRepo, datasystem.NewLogFileRepo, datasystem.NewTaskRepo,
datasystem.NewMediaRepo, datasystem.NewAnnouncementRepo, datapayment.NewPaymentRepo, datapayment.NewPaymentOrderRepo,
datasystem.NewIntegrationConfigRepo,
)
type Data struct {
initMu sync.Mutex
configMu sync.Mutex
databaseReady atomic.Bool
gormDB *reloadableDB
redis *reloadableRedis
mongo *reloadableMongo
runtime *conf.Runtime
storage *storage.Reloadable
dbListMu sync.RWMutex
dbList map[string]*gorm.DB
appLogger *slog.Logger
auditLog *dataScopeAuditWriter
catalog module.Catalog
}
// DB exposes the active primary database to narrowly scoped data submodules.
func (d *Data) DB() *gorm.DB {
if d == nil || d.gormDB == nil {
return nil
}
return d.gormDB.DB()
}
// DatabaseReady reports whether the configured primary database has been
// initialized. It is intentionally small so system repositories do not depend
// on the full Data implementation.
func (d *Data) DatabaseReady() bool {
return d != nil && d.databaseReady.Load()
}
// Runtime exposes the immutable runtime configuration snapshot to data
// submodules that need system settings while keeping Data itself private.
func (d *Data) Runtime() *conf.Runtime {
if d == nil {
return nil
}
return d.runtime
}
// Database resolves the primary or a named database for repositories such as
// the system export module.
func (d *Data) Database(name string) (*gorm.DB, error) {
return d.database(name)
}
func (d *Data) RedisClient() redis.UniversalClient {
if d == nil || d.redis == nil {
return nil
}
return d.redis.load()
}
func (d *Data) logger() *slog.Logger {
if d != nil && d.appLogger != nil {
return d.appLogger
}
return slog.Default()
}
func openDatabaseList(configs []*conf.Data_Database, appLogger ...*slog.Logger) (map[string]*gorm.DB, error) {
items := make(map[string]*gorm.DB)
for _, config := range configs {
if config == nil || config.Disable || config.AliasName == "" {
continue
}
db, err := openDatabase(config, false, "", appLogger...)
if err != nil {
for _, opened := range items {
if sqlDB, dbErr := opened.DB(); dbErr == nil {
_ = sqlDB.Close()
}
}
return nil, fmt.Errorf("open database %q: %w", config.AliasName, err)
}
items[config.AliasName] = db
}
return items, nil
}
func closeDatabaseList(items map[string]*gorm.DB) {
for _, db := range items {
if sqlDB, err := db.DB(); err == nil {
_ = sqlDB.Close()
}
}
}
func (d *Data) replaceDatabaseList(items map[string]*gorm.DB) {
d.dbListMu.Lock()
old := d.dbList
d.dbList = items
d.dbListMu.Unlock()
closeDatabaseList(old)
}
func (d *Data) database(name string) (*gorm.DB, error) {
if name == "" {
return d.gormDB.DB(), nil
}
d.dbListMu.RLock()
db := d.dbList[name]
d.dbListMu.RUnlock()
if db == nil {
return nil, fmt.Errorf("database %q not found", name)
}
return db, nil
}
func NewData(runtime *conf.Runtime, appLogger *slog.Logger, storageManager *storage.Reloadable, catalog module.Catalog) (*Data, func(), error) {
if appLogger == nil {
appLogger = slog.Default()
}
c := runtime.Data()
if c == nil {
c = &conf.Data{}
}
if c.Database == nil {
// An empty database block starts the service on the bootstrap database
// so /init/checkdb
// and /init/initdb remain available.
c.Database = &conf.Data_Database{}
}
d := &Data{runtime: runtime, appLogger: appLogger, storage: storageManager, catalog: catalog}
usingFallback := !databaseConnectionConfigured(c.Database)
var db *gorm.DB
var err error
if !usingFallback {
db, err = openDatabase(c.Database, false, "", appLogger)
}
if usingFallback || err != nil {
// The initialization endpoint must remain available when the configured
// target database has not been created yet.
if err != nil {
appLogger.Warn("configured database unavailable before initialization", "mod", "system", "error", err)
}
db, err = openFallbackDatabase(appLogger)
if err != nil {
return nil, nil, fmt.Errorf("open bootstrap database: %w", err)
}
usingFallback = true
}
d.databaseReady.Store(!usingFallback)
d.gormDB = newReloadableDB(db, d.enqueueDataScopeAudit)
d.auditLog = newDataScopeAuditWriter(d, appLogger)
d.dbList, err = openDatabaseList(c.DatabaseList, appLogger)
if err != nil {
d.auditLog.Close()
d.gormDB.close()
return nil, nil, err
}
for _, item := range d.dbList {
registerDataScopeCallbacks(item, d.enqueueDataScopeAudit)
}
admin := runtime.Admin()
if admin == nil {
admin = &conf.AdminBackend{}
}
disableAutoMigrate := admin.System != nil && admin.System.DisableAutoMigrate
if !usingFallback && !disableAutoMigrate {
if err = migrateAll(db, catalog); err != nil {
return nil, nil, fmt.Errorf("migrate tables: %w", err)
}
}
if !usingFallback {
storageConfig, storageErr := resolveStorageIntegrationConfig(db, admin.Storage)
if storageErr != nil {
return nil, nil, fmt.Errorf("load storage integration configuration: %w", storageErr)
}
emailConfig, emailErr := resolveEmailIntegrationConfig(db, admin.Email)
if emailErr != nil {
return nil, nil, fmt.Errorf("load email integration configuration: %w", emailErr)
}
admin.Storage = storageConfig
admin.Email = emailConfig
websocketConfig, websocketErr := resolveWebSocketIntegrationConfig(db, admin.Websocket)
if websocketErr != nil {
return nil, nil, fmt.Errorf("load websocket integration configuration: %w", websocketErr)
}
admin.Websocket = websocketConfig
mqConfig, mqErr := resolveMQIntegrationConfig(db, admin.Mq)
if mqErr != nil {
return nil, nil, fmt.Errorf("load mq integration configuration: %w", mqErr)
}
admin.Mq = mqConfig
runtime.Replace(c, admin)
activeStorage, storageErr := storage.New(admin)
if storageErr != nil {
return nil, nil, fmt.Errorf("initialize storage: %w", storageErr)
}
if storageManager != nil {
storageManager.Replace(activeStorage)
}
if db.Migrator().HasTable(&integrationConfigPO{}) {
if removeErr := d.removeIntegrationConfigFromFile(); removeErr != nil {
appLogger.Warn("remove legacy integration configuration from file", "mod", "integration", "error", removeErr)
}
}
}
if storageManager == nil {
return nil, nil, fmt.Errorf("storage manager is nil")
}
useRedis := admin != nil && admin.System != nil && admin.System.UseRedis
d.redis = newReloadableRedis(openRedis(c.Redis, useRedis, appLogger))
useMongo := admin != nil && admin.System != nil && admin.System.UseMongo
mongoClient, err := openMongo(c.Mongo, useMongo)
if err != nil {
appLogger.Error("mongo unavailable", "mod", "mongo", "error", err)
mongoClient = nil
}
d.mongo = newReloadableMongo(mongoClient)
stopConfigWatcher := d.watchConfig()
cleanup := func() {
stopConfigWatcher()
d.auditLog.Close()
d.gormDB.close()
closeDatabaseList(d.dbList)
d.redis.close()
d.mongo.close()
}
return d, cleanup, nil
}
func openRedis(config *conf.Data_Redis, enabled bool, appLogger ...*slog.Logger) redis.UniversalClient {
if !enabled || config == nil || (config.Addr == "" && len(config.ClusterAddrs) == 0) {
return nil
}
var candidate redis.UniversalClient
if config.UseCluster {
addresses := config.ClusterAddrs
if len(addresses) == 0 && config.Addr != "" {
addresses = []string{config.Addr}
}
candidate = redis.NewClusterClient(&redis.ClusterOptions{Addrs: addresses, Password: config.Password})
} else {
options := &redis.Options{Addr: config.Addr, Network: config.Network, Password: config.Password, DB: int(config.Db)}
if config.ReadTimeout != nil {
options.ReadTimeout = config.ReadTimeout.AsDuration()
}
if config.WriteTimeout != nil {
options.WriteTimeout = config.WriteTimeout.AsDuration()
}
candidate = redis.NewClient(options)
}
pingCtx, cancel := context.WithTimeout(context.Background(), 800*time.Millisecond)
defer cancel()
if err := candidate.Ping(pingCtx).Err(); err != nil {
log := slog.Default()
if len(appLogger) > 0 && appLogger[0] != nil {
log = appLogger[0]
}
log.Warn("redis unavailable, using in-memory cache", "mod", "redis", "error", err)
_ = candidate.Close()
return nil
}
return candidate
}
func (d *Data) activateDatabase(db *gorm.DB, config *conf.Data_Database) {
d.gormDB.replace(db, d.enqueueDataScopeAudit)
d.runtime.UpdateDatabase(config)
d.databaseReady.Store(true)
}