You can not select more than 25 topics
Topics must start with a letter or number, can include dashes ('-') and can be up to 35 characters long.
166 lines
5.1 KiB
166 lines
5.1 KiB
package mysql
|
|
|
|
import (
|
|
"errors"
|
|
"fmt"
|
|
"time"
|
|
|
|
mysqldrv "github.com/go-sql-driver/mysql"
|
|
"gorm.io/driver/mysql"
|
|
"gorm.io/gorm"
|
|
)
|
|
|
|
func newSys(options Options) (sys *MySql, err error) {
|
|
sys = &MySql{options: options}
|
|
err = sys.init()
|
|
return
|
|
}
|
|
|
|
type MySql struct {
|
|
options Options
|
|
db *gorm.DB
|
|
}
|
|
|
|
func (this *MySql) init() (err error) {
|
|
this.db, err = gorm.Open(mysql.Open(this.options.MySQLDsn), &gorm.Config{})
|
|
if err != nil {
|
|
this.options.Log.Errorln(err)
|
|
}
|
|
return
|
|
}
|
|
|
|
// AutoIncrementStart 设置表自增主键的起始值(MySQL: ALTER TABLE ... AUTO_INCREMENT)。
|
|
// column 仅为跨驱动接口对齐保留,MySQL 用表级设置无需列名。
|
|
func (this *MySql) AutoIncrementStart(tName, column string, start uint64) (err error) {
|
|
_ = column
|
|
err = this.db.Exec(fmt.Sprintf("ALTER TABLE %s AUTO_INCREMENT = %d", tName, start)).Error
|
|
return
|
|
}
|
|
|
|
// AutoIncrementFloor 抬高自增主键下限:新插入的行 id >= start,已有行不动。
|
|
// MySQL 的 ALTER TABLE ... AUTO_INCREMENT 本身就是下限语义(低于当前最大值+1 时被静默忽略),
|
|
// 所以与 AutoIncrementStart 同一条语句;两个方法并存只为与 postgres 驱动的接口对齐。
|
|
func (this *MySql) AutoIncrementFloor(tName, column string, start uint64) (err error) {
|
|
return this.AutoIncrementStart(tName, column, start)
|
|
}
|
|
|
|
func (this *MySql) Exec(sql string, values ...interface{}) (tx *gorm.DB) {
|
|
tx = this.db.Exec(sql, values...)
|
|
return
|
|
}
|
|
|
|
func (this *MySql) Raw(sql string, values ...interface{}) (tx *gorm.DB) {
|
|
tx = this.db.Raw(sql, values...)
|
|
return
|
|
}
|
|
|
|
// CreateTable 建表 / 同步表结构(gorm AutoMigrate)。
|
|
//
|
|
// ⚠️ 同一张表可能被多个进程在启动时同时同步:业务容器里 gateway/home/api/mcp/timer 是一起拉起的,
|
|
// 例如 echomeet_record 就同时由 home 的 echomeet 模块与 api 模块建。AutoMigrate 是「先查列在不在、
|
|
// 再 ALTER」两步,不是原子的——两边都查到缺列、都去加,后到的那个必然报 1060 Duplicate column,
|
|
// 调用方把它当初始化失败,整个模块 panic、接口全部 code:11。2026-09-26 正式服发 0.2.9 时
|
|
// home 就这样挂了几分钟(新加 computeuser/computedevices 两列,api 抢先加成功)。
|
|
// 所以遇到这类「对方已经建好了」的冲突,等一下再同步一次:第二遍能看到对方加好的列,自然通过。
|
|
func (this *MySql) CreateTable(tName string, model any) (err error) {
|
|
for attempt := 0; attempt < 3; attempt++ {
|
|
if attempt > 0 {
|
|
time.Sleep(time.Duration(attempt) * 500 * time.Millisecond)
|
|
}
|
|
if err = this.db.Table(tName).AutoMigrate(model); err == nil || !isConcurrentDDLConflict(err) {
|
|
return
|
|
}
|
|
this.options.Log.Warnf("mysql CreateTable %s 与其它进程并发建表冲突,稍后重试(第%d次): %v", tName, attempt+1, err)
|
|
}
|
|
return
|
|
}
|
|
|
|
// isConcurrentDDLConflict 是否是「别的进程已经把这个表/列/索引建好了」这一类冲突:
|
|
// 1050 表已存在、1060 列已存在、1061 索引名已存在。
|
|
func isConcurrentDDLConflict(err error) bool {
|
|
var me *mysqldrv.MySQLError
|
|
if !errors.As(err, &me) {
|
|
return false
|
|
}
|
|
switch me.Number {
|
|
case 1050, 1060, 1061:
|
|
return true
|
|
}
|
|
return false
|
|
}
|
|
|
|
// 获取表对象
|
|
func (this *MySql) Table(tName string) (tx *gorm.DB) {
|
|
tx = this.db.Table(tName)
|
|
return
|
|
}
|
|
|
|
// 查询数据
|
|
func (this *MySql) FindOne(tName string, model any, query interface{}, args ...interface{}) (err error) {
|
|
result := this.db.Table(tName).Where(query, args...).First(model)
|
|
if result.Error != nil {
|
|
err = result.Error
|
|
}
|
|
return
|
|
}
|
|
|
|
// FindOnePrimary MySQL 未做读写分离,等同 FindOne(仅为与 postgres 驱动接口对齐)。
|
|
func (this *MySql) FindOnePrimary(tName string, model any, query interface{}, args ...interface{}) (err error) {
|
|
return this.FindOne(tName, model, query, args...)
|
|
}
|
|
|
|
// 查询数据
|
|
func (this *MySql) Find(tName string, models any, query interface{}, args ...interface{}) (err error) {
|
|
result := this.db.Table(tName).Where(query, args...).Find(models)
|
|
if result.Error != nil {
|
|
err = result.Error
|
|
}
|
|
return
|
|
}
|
|
|
|
// 插入数据
|
|
func (this *MySql) Insert(tName string, model any) (err error) {
|
|
result := this.db.Table(tName).Create(model)
|
|
if result.Error != nil {
|
|
err = result.Error
|
|
}
|
|
return
|
|
}
|
|
|
|
// 插入数据
|
|
func (this *MySql) Update(tName string, where, change map[string]interface{}) (err error) {
|
|
// 使用 Update 方法更新特定字段
|
|
result := this.db.Table(tName).Where(where).Updates(change)
|
|
if result.Error != nil {
|
|
err = result.Error
|
|
}
|
|
return
|
|
}
|
|
|
|
func (this *MySql) Save(tName string, model any) (err error) {
|
|
err = this.db.Table(tName).Save(model).Error
|
|
return
|
|
}
|
|
|
|
// 插入数据
|
|
func (this *MySql) Delete(tName string, query interface{}, args ...interface{}) (err error) {
|
|
// 使用 Update 方法更新特定字段
|
|
result := this.db.Table(tName).Where(query, args...).Unscoped().Delete(nil)
|
|
if result.Error != nil {
|
|
err = result.Error
|
|
}
|
|
return
|
|
}
|
|
|
|
// 删除表
|
|
func (this *MySql) DropTable(tName string) (err error) {
|
|
// 使用 Update 方法更新特定字段
|
|
err = this.db.Migrator().DropTable(tName)
|
|
return
|
|
}
|
|
|
|
// 事务
|
|
func (this *MySql) Begin() (tx *gorm.DB) {
|
|
tx = this.db.Begin()
|
|
return
|
|
}
|
|
|