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) ensureUserPasswordHashColumn(db) ensureSuperArticlesTableRaw(db) ensureSuperArticleImagesColumn(db) ensureSuperArticleAuditColumns(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.Book{}); err != nil { log.Printf("database: books 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.SuperArticle{}); err != nil { log.Printf("database: super_articles 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) ensureUserPasswordHashColumn(db) ensureSuperArticlesTableRaw(db) ensureSuperArticleImagesColumn(db) ensureSuperArticleAuditColumns(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) } } } // ensureUserPasswordHashColumn users 表补充 password_hash(SKIP_AUTO_MIGRATE 等场景) func ensureUserPasswordHashColumn(db *gorm.DB) { m := db.Migrator() if !m.HasColumn(&model.User{}, "PasswordHash") { if err := db.Exec("ALTER TABLE users ADD COLUMN password_hash VARCHAR(128) NULL COMMENT 'H5登录 bcrypt'").Error; err != nil { log.Printf("database: users add password_hash 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) } } } // ensureSuperArticlesTableRaw SKIP_AUTO_MIGRATE 或未跑 AutoMigrate 时兜底建表(与 model.SuperArticle、scripts/create_super_articles.sql 一致) func ensureSuperArticlesTableRaw(db *gorm.DB) { m := db.Migrator() if m.HasTable(&model.SuperArticle{}) { return } sqlStmt := ` CREATE TABLE IF NOT EXISTS super_articles ( id BIGINT UNSIGNED NOT NULL AUTO_INCREMENT, user_id VARCHAR(50) NOT NULL COMMENT '作者 users.id', title VARCHAR(200) NOT NULL DEFAULT '', content LONGTEXT COMMENT '正文', images LONGTEXT NULL COMMENT '配图 URL JSON 数组', audit_status VARCHAR(32) NOT NULL DEFAULT 'pending' COMMENT 'pending/approved/rejected', reject_reason LONGTEXT NULL COMMENT '驳回原因', created_at DATETIME(3) NULL, updated_at DATETIME(3) NULL, PRIMARY KEY (id), KEY idx_super_articles_user_time (user_id, created_at), KEY idx_super_articles_audit_status (audit_status) ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_unicode_ci COMMENT='超级个体发文(小程序 /s/:id H5)'` if err := db.Exec(sqlStmt).Error; err != nil { log.Printf("database: ensureSuperArticlesTableRaw warning: %v", err) } } // ensureSuperArticleImagesColumn 旧库补 images 列(SKIP_AUTO_MIGRATE / 早期建表无该列) func ensureSuperArticleImagesColumn(db *gorm.DB) { m := db.Migrator() if !m.HasTable(&model.SuperArticle{}) { return } if m.HasColumn(&model.SuperArticle{}, "Images") { return } if err := db.Exec("ALTER TABLE super_articles ADD COLUMN images LONGTEXT NULL COMMENT '配图 URL JSON 数组'").Error; err != nil { msg := strings.ToLower(err.Error()) if strings.Contains(msg, "duplicate column") { return } log.Printf("database: super_articles add images column warning: %v", err) } } // ensureSuperArticleAuditColumns 旧库补审核字段;历史数据一律视为已通过,避免线上动态消失 func ensureSuperArticleAuditColumns(db *gorm.DB) { m := db.Migrator() if !m.HasTable(&model.SuperArticle{}) { return } if !m.HasColumn(&model.SuperArticle{}, "AuditStatus") { if err := db.Exec("ALTER TABLE super_articles ADD COLUMN audit_status VARCHAR(32) NOT NULL DEFAULT 'approved' COMMENT 'pending/approved/rejected'").Error; err != nil { msg := strings.ToLower(err.Error()) if !strings.Contains(msg, "duplicate column") { log.Printf("database: super_articles add audit_status warning: %v", err) } } } if !m.HasColumn(&model.SuperArticle{}, "RejectReason") { if err := db.Exec("ALTER TABLE super_articles ADD COLUMN reject_reason LONGTEXT NULL COMMENT '驳回原因'").Error; err != nil { msg := strings.ToLower(err.Error()) if !strings.Contains(msg, "duplicate column") { log.Printf("database: super_articles add reject_reason warning: %v", err) } } } _ = db.Exec("UPDATE super_articles SET audit_status = 'approved' WHERE audit_status IS NULL OR TRIM(audit_status) = ''").Error if !m.HasIndex(&model.SuperArticle{}, "idx_super_articles_audit_status") { if err := db.Exec("CREATE INDEX idx_super_articles_audit_status ON super_articles(audit_status)").Error; err != nil { msg := strings.ToLower(err.Error()) if !strings.Contains(msg, "duplicate") { log.Printf("database: super_articles idx_super_articles_audit_status 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) } } }