320 lines
14 KiB
Go
320 lines
14 KiB
Go
package database
|
||
|
||
import (
|
||
"log"
|
||
"os"
|
||
"strconv"
|
||
"strings"
|
||
"time"
|
||
|
||
"soul-api/internal/model"
|
||
|
||
"gorm.io/driver/mysql"
|
||
"gorm.io/gorm"
|
||
"gorm.io/gorm/logger"
|
||
)
|
||
|
||
var db *gorm.DB
|
||
|
||
// Init 使用 DSN 连接 MySQL,供 handler 通过 DB() 使用
|
||
func Init(dsn string) error {
|
||
// 慢查询阈值:默认 5 秒,避免 GORM 默认 200ms 导致控制台刷屏;可通过 SLOW_SQL_THRESHOLD_MS 覆盖
|
||
slowMs := 5000
|
||
if s := os.Getenv("SLOW_SQL_THRESHOLD_MS"); s != "" {
|
||
if n, e := strconv.Atoi(s); e == nil && n > 0 {
|
||
slowMs = n
|
||
}
|
||
}
|
||
gormLogger := logger.New(
|
||
log.New(os.Stdout, "\r\n", log.LstdFlags),
|
||
logger.Config{
|
||
SlowThreshold: time.Duration(slowMs) * time.Millisecond,
|
||
IgnoreRecordNotFoundError: true,
|
||
Colorful: true,
|
||
},
|
||
)
|
||
|
||
var err error
|
||
db, err = gorm.Open(mysql.Open(dsn), &gorm.Config{Logger: gormLogger})
|
||
if err != nil {
|
||
return err
|
||
}
|
||
|
||
// 连接池配置:防止腾讯云 CDB 等云数据库在空闲后主动断开长连接导致 "unexpected EOF"
|
||
sqlDB, err := db.DB()
|
||
if err != nil {
|
||
return err
|
||
}
|
||
sqlDB.SetMaxOpenConns(25)
|
||
sqlDB.SetMaxIdleConns(5)
|
||
sqlDB.SetConnMaxLifetime(5 * time.Minute) // 短于 MySQL wait_timeout,避免用到已关闭的连接
|
||
sqlDB.SetConnMaxIdleTime(2 * time.Minute) // 空闲连接超过此时间自动关闭并重建
|
||
|
||
skipMigrate := strings.ToLower(strings.TrimSpace(os.Getenv("SKIP_AUTO_MIGRATE")))
|
||
if skipMigrate == "1" || skipMigrate == "true" || skipMigrate == "yes" {
|
||
log.Println("database: SKIP_AUTO_MIGRATE enabled, skipping schema migration")
|
||
// 即使跳过 AutoMigrate,也补齐关键运行时字段,避免新功能因历史库缺列直接报错。
|
||
ensurePersonSchema(db)
|
||
ensureCkbLeadSchema(db)
|
||
ensureLinkTagSchema(db)
|
||
// 开放平台 API Key / 请求日志表:全量 AutoMigrate 被跳过时仍需建表,否则管理端报 1146
|
||
if err := db.AutoMigrate(&model.OpenPlatformApiKey{}, &model.OpenPlatformApiLog{}); err != nil {
|
||
log.Printf("database: open_platform tables migrate warning: %v", err)
|
||
}
|
||
ensureOpenPlatformTablesRaw(db)
|
||
ensureUserDiscPdpColumns(db)
|
||
log.Println("database: connected")
|
||
return nil
|
||
}
|
||
|
||
if err := db.AutoMigrate(&model.WechatCallbackLog{}); err != nil {
|
||
log.Printf("database: wechat_callback_logs migrate warning: %v", err)
|
||
}
|
||
if err := db.AutoMigrate(&model.Withdrawal{}); err != nil {
|
||
log.Printf("database: withdrawals migrate warning: %v", err)
|
||
}
|
||
if err := db.AutoMigrate(&model.MatchRecord{}); err != nil {
|
||
log.Printf("database: match_records migrate warning: %v", err)
|
||
}
|
||
if err := db.AutoMigrate(&model.UserAddress{}); err != nil {
|
||
log.Printf("database: user_addresses migrate warning: %v", err)
|
||
}
|
||
if err := db.AutoMigrate(&model.VipRole{}); err != nil {
|
||
log.Printf("database: vip_roles migrate warning: %v", err)
|
||
}
|
||
if err := db.AutoMigrate(&model.Order{}); err != nil {
|
||
log.Printf("database: orders migrate warning: %v", err)
|
||
}
|
||
if err := db.AutoMigrate(&model.UserBalance{}); err != nil {
|
||
log.Printf("database: user_balances migrate warning: %v", err)
|
||
}
|
||
if err := db.AutoMigrate(&model.BalanceTransaction{}); err != nil {
|
||
log.Printf("database: balance_transactions migrate warning: %v", err)
|
||
}
|
||
if err := db.AutoMigrate(&model.Mentor{}); err != nil {
|
||
log.Printf("database: mentors migrate warning: %v", err)
|
||
}
|
||
if err := db.AutoMigrate(&model.MentorConsultation{}); err != nil {
|
||
log.Printf("database: mentor_consultations migrate warning: %v", err)
|
||
}
|
||
if err := db.AutoMigrate(&model.AuthorConfig{}); err != nil {
|
||
log.Printf("database: author_config migrate warning: %v", err)
|
||
}
|
||
if err := db.AutoMigrate(&model.AdminUser{}); err != nil {
|
||
log.Printf("database: admin_users migrate warning: %v", err)
|
||
}
|
||
if err := db.AutoMigrate(&model.CkbLeadRecord{}); err != nil {
|
||
log.Printf("database: ckb_lead_records migrate warning: %v", err)
|
||
}
|
||
ensureCkbLeadSchema(db)
|
||
if err := db.AutoMigrate(&model.Person{}); err != nil {
|
||
log.Printf("database: persons migrate warning: %v", err)
|
||
}
|
||
// persons 历史库可能因旧索引冲突导致 AutoMigrate 中断,补一层列级自愈,避免 /api/db/persons 报 Unknown column。
|
||
ensurePersonSchema(db)
|
||
if err := db.AutoMigrate(&model.LinkTag{}); err != nil {
|
||
log.Printf("database: link_tags migrate warning: %v", err)
|
||
}
|
||
ensureLinkTagSchema(db)
|
||
// 以下表业务大量使用,必须参与 AutoMigrate,否则旧库缺字段会导致订单/用户/VIP 等接口报错
|
||
if err := db.AutoMigrate(&model.User{}); err != nil {
|
||
log.Printf("database: users migrate warning: %v", err)
|
||
}
|
||
if err := db.AutoMigrate(&model.SystemConfig{}); err != nil {
|
||
log.Printf("database: system_config migrate warning: %v", err)
|
||
}
|
||
if err := db.AutoMigrate(&model.Chapter{}); err != nil {
|
||
log.Printf("database: chapters migrate warning: %v", err)
|
||
}
|
||
if err := db.AutoMigrate(&model.GiftPayRequest{}); err != nil {
|
||
log.Printf("database: gift_pay_requests migrate warning: %v", err)
|
||
}
|
||
if err := db.AutoMigrate(&model.UserRule{}); err != nil {
|
||
log.Printf("database: user_rules migrate warning: %v", err)
|
||
}
|
||
if err := db.AutoMigrate(&model.UserTrack{}); err != nil {
|
||
log.Printf("database: user_tracks migrate warning: %v", err)
|
||
}
|
||
if err := db.AutoMigrate(&model.UserRuleCompletion{}); err != nil {
|
||
log.Printf("database: user_rule_completions migrate warning: %v", err)
|
||
}
|
||
if err := db.AutoMigrate(&model.OpenPlatformApiKey{}); err != nil {
|
||
log.Printf("database: open_platform_api_keys migrate warning: %v", err)
|
||
}
|
||
if err := db.AutoMigrate(&model.OpenPlatformApiLog{}); err != nil {
|
||
log.Printf("database: open_platform_api_logs migrate warning: %v", err)
|
||
}
|
||
ensureOpenPlatformTablesRaw(db)
|
||
ensureUserDiscPdpColumns(db)
|
||
log.Println("database: connected")
|
||
return nil
|
||
}
|
||
|
||
// ensureUserDiscPdpColumns users 表补充 disc/pdp(SKIP_AUTO_MIGRATE 时 User 未跑 AutoMigrate 时需手工列)
|
||
func ensureUserDiscPdpColumns(db *gorm.DB) {
|
||
m := db.Migrator()
|
||
if !m.HasColumn(&model.User{}, "disc") {
|
||
if err := db.Exec("ALTER TABLE users ADD COLUMN disc VARCHAR(64) NULL COMMENT 'DISC 测评'").Error; err != nil {
|
||
log.Printf("database: users add disc column warning: %v", err)
|
||
}
|
||
}
|
||
if !m.HasColumn(&model.User{}, "pdp") {
|
||
if err := db.Exec("ALTER TABLE users ADD COLUMN pdp VARCHAR(64) NULL COMMENT 'PDP 测评'").Error; err != nil {
|
||
log.Printf("database: users add pdp column warning: %v", err)
|
||
}
|
||
}
|
||
}
|
||
|
||
// DB 返回全局 *gorm.DB,仅在 Init 成功后调用
|
||
func DB() *gorm.DB {
|
||
return db
|
||
}
|
||
|
||
// ensureOpenPlatformTablesRaw AutoMigrate 失败或旧库无表时,用 CREATE IF NOT EXISTS 兜底(与 model 字段一致)
|
||
func ensureOpenPlatformTablesRaw(db *gorm.DB) {
|
||
m := db.Migrator()
|
||
if m.HasTable(&model.OpenPlatformApiKey{}) && m.HasTable(&model.OpenPlatformApiLog{}) {
|
||
return
|
||
}
|
||
sqlKeys := `
|
||
CREATE TABLE IF NOT EXISTS open_platform_api_keys (
|
||
id BIGINT UNSIGNED NOT NULL AUTO_INCREMENT,
|
||
name VARCHAR(120) NOT NULL DEFAULT '',
|
||
key_prefix VARCHAR(48) NOT NULL,
|
||
key_hash VARCHAR(64) NOT NULL,
|
||
revoked_at DATETIME(3) NULL,
|
||
created_at DATETIME(3) NULL,
|
||
updated_at DATETIME(3) NULL,
|
||
PRIMARY KEY (id),
|
||
KEY idx_open_platform_api_keys_revoked_at (revoked_at)
|
||
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_unicode_ci`
|
||
sqlLogs := `
|
||
CREATE TABLE IF NOT EXISTS open_platform_api_logs (
|
||
id BIGINT UNSIGNED NOT NULL AUTO_INCREMENT,
|
||
api_key_id BIGINT UNSIGNED NOT NULL,
|
||
key_prefix VARCHAR(48) NOT NULL DEFAULT '',
|
||
method VARCHAR(16) NOT NULL DEFAULT '',
|
||
path VARCHAR(512) NOT NULL DEFAULT '',
|
||
status_code INT NOT NULL DEFAULT 0,
|
||
client_ip VARCHAR(64) NOT NULL DEFAULT '',
|
||
duration_ms INT NOT NULL DEFAULT 0,
|
||
request_body TEXT NULL,
|
||
response_body TEXT NULL,
|
||
created_at DATETIME(3) NULL,
|
||
PRIMARY KEY (id),
|
||
KEY idx_open_platform_api_logs_api_key_id (api_key_id),
|
||
KEY idx_open_platform_api_logs_created_at (created_at)
|
||
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_unicode_ci`
|
||
if !m.HasTable(&model.OpenPlatformApiKey{}) {
|
||
if err := db.Exec(sqlKeys).Error; err != nil {
|
||
log.Printf("database: ensureOpenPlatformTablesRaw open_platform_api_keys warning: %v", err)
|
||
}
|
||
}
|
||
if !m.HasTable(&model.OpenPlatformApiLog{}) {
|
||
if err := db.Exec(sqlLogs).Error; err != nil {
|
||
log.Printf("database: ensureOpenPlatformTablesRaw open_platform_api_logs warning: %v", err)
|
||
}
|
||
}
|
||
}
|
||
|
||
func ensurePersonSchema(db *gorm.DB) {
|
||
m := db.Migrator()
|
||
if !m.HasColumn(&model.Person{}, "is_pinned") {
|
||
if err := db.Exec("ALTER TABLE persons ADD COLUMN is_pinned TINYINT(1) NOT NULL DEFAULT 0 COMMENT '置顶到小程序首页'").Error; err != nil {
|
||
log.Printf("database: persons schema ensure warning: %v; action=add is_pinned", err)
|
||
}
|
||
}
|
||
if !m.HasColumn(&model.Person{}, "person_source") {
|
||
if err := db.Exec("ALTER TABLE persons ADD COLUMN person_source VARCHAR(32) NOT NULL DEFAULT '' COMMENT '来源:空=后台手工;vip_sync=超级个体同步'").Error; err != nil {
|
||
log.Printf("database: persons schema ensure warning: %v; action=add person_source", err)
|
||
}
|
||
}
|
||
if !m.HasIndex(&model.Person{}, "idx_persons_is_pinned") {
|
||
if err := db.Exec("CREATE INDEX idx_persons_is_pinned ON persons(is_pinned)").Error; err != nil {
|
||
log.Printf("database: persons schema ensure warning: %v; action=create idx_persons_is_pinned", err)
|
||
}
|
||
}
|
||
if !m.HasColumn(&model.Person{}, "home_entry_config") {
|
||
if err := db.Exec("ALTER TABLE persons ADD COLUMN home_entry_config TEXT NULL COMMENT '首页入口配置(打赏/上麦,JSON)'").Error; err != nil {
|
||
log.Printf("database: persons schema ensure warning: %v; action=add home_entry_config", err)
|
||
}
|
||
}
|
||
}
|
||
|
||
func ensureCkbLeadSchema(db *gorm.DB) {
|
||
m := db.Migrator()
|
||
if !m.HasColumn(&model.CkbLeadRecord{}, "push_status") {
|
||
if err := db.Exec("ALTER TABLE ckb_lead_records ADD COLUMN push_status VARCHAR(20) NOT NULL DEFAULT 'pending' COMMENT '推送状态: pending/success/failed'").Error; err != nil {
|
||
log.Printf("database: ckb_lead_records schema ensure warning: %v; action=add push_status", err)
|
||
}
|
||
}
|
||
if !m.HasColumn(&model.CkbLeadRecord{}, "action") {
|
||
if err := db.Exec("ALTER TABLE ckb_lead_records ADD COLUMN action VARCHAR(20) NOT NULL DEFAULT 'lead' COMMENT '记录类型: lead/join/match'").Error; err != nil {
|
||
log.Printf("database: ckb_lead_records schema ensure warning: %v; action=add action", err)
|
||
}
|
||
}
|
||
if !m.HasColumn(&model.CkbLeadRecord{}, "plan_api_key") {
|
||
if err := db.Exec("ALTER TABLE ckb_lead_records ADD COLUMN plan_api_key VARCHAR(100) NOT NULL DEFAULT '' COMMENT '本次命中的获客计划 apiKey 快照'").Error; err != nil {
|
||
log.Printf("database: ckb_lead_records schema ensure warning: %v; action=add plan_api_key", err)
|
||
}
|
||
}
|
||
if !m.HasIndex(&model.CkbLeadRecord{}, "idx_ckb_lead_action") {
|
||
if err := db.Exec("CREATE INDEX idx_ckb_lead_action ON ckb_lead_records(action)").Error; err != nil {
|
||
log.Printf("database: ckb_lead_records schema ensure warning: %v; action=create idx_ckb_lead_action", err)
|
||
}
|
||
}
|
||
if !m.HasColumn(&model.CkbLeadRecord{}, "ckb_code") {
|
||
if err := db.Exec("ALTER TABLE ckb_lead_records ADD COLUMN ckb_code INT NOT NULL DEFAULT 0 COMMENT '存客宝响应 code(快照)'").Error; err != nil {
|
||
log.Printf("database: ckb_lead_records schema ensure warning: %v; action=add ckb_code", err)
|
||
}
|
||
}
|
||
if !m.HasColumn(&model.CkbLeadRecord{}, "ckb_message") {
|
||
if err := db.Exec("ALTER TABLE ckb_lead_records ADD COLUMN ckb_message VARCHAR(500) NOT NULL DEFAULT '' COMMENT '存客宝响应 message(快照)'").Error; err != nil {
|
||
log.Printf("database: ckb_lead_records schema ensure warning: %v; action=add ckb_message", err)
|
||
}
|
||
}
|
||
if !m.HasColumn(&model.CkbLeadRecord{}, "ckb_data") {
|
||
if err := db.Exec("ALTER TABLE ckb_lead_records ADD COLUMN ckb_data TEXT NULL COMMENT '存客宝响应 data(快照 JSON)'").Error; err != nil {
|
||
log.Printf("database: ckb_lead_records schema ensure warning: %v; action=add ckb_data", err)
|
||
}
|
||
}
|
||
if !m.HasColumn(&model.CkbLeadRecord{}, "retry_count") {
|
||
if err := db.Exec("ALTER TABLE ckb_lead_records ADD COLUMN retry_count INT NOT NULL DEFAULT 0 COMMENT '重试次数'").Error; err != nil {
|
||
log.Printf("database: ckb_lead_records schema ensure warning: %v; action=add retry_count", err)
|
||
}
|
||
}
|
||
if !m.HasColumn(&model.CkbLeadRecord{}, "last_push_at") {
|
||
if err := db.Exec("ALTER TABLE ckb_lead_records ADD COLUMN last_push_at DATETIME NULL COMMENT '最后推送时间'").Error; err != nil {
|
||
log.Printf("database: ckb_lead_records schema ensure warning: %v; action=add last_push_at", err)
|
||
}
|
||
}
|
||
if !m.HasColumn(&model.CkbLeadRecord{}, "next_retry_at") {
|
||
if err := db.Exec("ALTER TABLE ckb_lead_records ADD COLUMN next_retry_at DATETIME NULL COMMENT '下次重试时间'").Error; err != nil {
|
||
log.Printf("database: ckb_lead_records schema ensure warning: %v; action=add next_retry_at", err)
|
||
}
|
||
}
|
||
if !m.HasIndex(&model.CkbLeadRecord{}, "idx_ckb_lead_push_status") {
|
||
if err := db.Exec("CREATE INDEX idx_ckb_lead_push_status ON ckb_lead_records(push_status)").Error; err != nil {
|
||
log.Printf("database: ckb_lead_records schema ensure warning: %v; action=create idx_ckb_lead_push_status", err)
|
||
}
|
||
}
|
||
// 放宽 push_status 长度以容纳存客宝细粒度状态(pending_verify、expired 等)
|
||
if err := db.Exec("ALTER TABLE ckb_lead_records MODIFY COLUMN push_status VARCHAR(64) NOT NULL DEFAULT 'pending' COMMENT '推送状态'").Error; err != nil {
|
||
log.Printf("database: ckb_lead_records schema ensure warning: %v; action=widen push_status", err)
|
||
}
|
||
}
|
||
|
||
func ensureLinkTagSchema(db *gorm.DB) {
|
||
m := db.Migrator()
|
||
if !m.HasColumn(&model.LinkTag{}, "pass_phone") {
|
||
if err := db.Exec("ALTER TABLE link_tags ADD COLUMN pass_phone TINYINT(1) NOT NULL DEFAULT 0 COMMENT '是否透传手机号'").Error; err != nil {
|
||
log.Printf("database: link_tags schema ensure warning: %v; action=add pass_phone", err)
|
||
}
|
||
}
|
||
if !m.HasColumn(&model.LinkTag{}, "phone_param_name") {
|
||
if err := db.Exec("ALTER TABLE link_tags ADD COLUMN phone_param_name VARCHAR(64) NOT NULL DEFAULT 'phone' COMMENT '手机号透传参数名'").Error; err != nil {
|
||
log.Printf("database: link_tags schema ensure warning: %v; action=add phone_param_name", err)
|
||
}
|
||
}
|
||
}
|