Merge branch 'master' into main-dmdb

# Conflicts:
#	db/db.go
This commit is contained in:
2026-04-11 21:54:37 +08:00
3 changed files with 179 additions and 10 deletions
+14 -1
View File
@@ -95,7 +95,12 @@ func (that *HoTimeDB) queryWithRetry(query string, retried bool, args ...interfa
that.LastErr.SetError(err)
if err != nil {
if !retried {
if that.Tx != nil {
if that.txFailed != nil {
*that.txFailed = true
}
}
if !retried && that.Tx == nil {
if pingErr := db.Ping(); pingErr == nil {
return that.queryWithRetry(query, true, args...)
}
@@ -169,6 +174,14 @@ func (that *HoTimeDB) execWithRetry(query string, retried bool, args ...interfac
if that.testTx != nil {
return resl, that.LastErr
}
// 事务内不做连接级重试:死锁等错误会导致 MySQL 自动回滚事务,
// 在已回滚的 Tx 上重试会以 auto-commit 模式执行,造成数据不一致
if that.Tx != nil {
if that.txFailed != nil {
*that.txFailed = true
}
return resl, that.LastErr
}
if !retried {
if pingErr := that.DB.Ping(); pingErr == nil {
return that.execWithRetry(query, true, args...)
+30 -9
View File
@@ -39,21 +39,36 @@ func (that *HoTimeDB) Action(action func(db HoTimeDB) (isSuccess bool)) (isSucce
testLogBuf: logBuf,
}
txFailed := false
db.txFailed = &txFailed
if that.testTx != nil {
spName := fmt.Sprintf("sp_%d", atomic.AddUint64(&savepointCounter, 1))
if that.testMu != nil { that.testMu.Lock() }
if that.testMu != nil {
that.testMu.Lock()
}
_, _ = that.testTx.Exec("SAVEPOINT " + spName)
if that.testMu != nil { that.testMu.Unlock() }
if that.testMu != nil {
that.testMu.Unlock()
}
db.Tx = that.testTx
isSuccess = action(db)
if !isSuccess {
if that.testMu != nil { that.testMu.Lock() }
if txFailed || !isSuccess {
if that.testMu != nil {
that.testMu.Lock()
}
_, _ = that.testTx.Exec("ROLLBACK TO SAVEPOINT " + spName)
if that.testMu != nil { that.testMu.Unlock() }
} else {
if that.testMu != nil { that.testMu.Lock() }
_, _ = that.testTx.Exec("RELEASE SAVEPOINT " + spName)
if that.testMu != nil { that.testMu.Unlock() }
if that.testMu != nil {
that.testMu.Unlock()
}
return false
}
if that.testMu != nil {
that.testMu.Lock()
}
_, _ = that.testTx.Exec("RELEASE SAVEPOINT " + spName)
if that.testMu != nil {
that.testMu.Unlock()
}
return isSuccess
}
@@ -69,6 +84,12 @@ func (that *HoTimeDB) Action(action func(db HoTimeDB) (isSuccess bool)) (isSucce
isSuccess = action(db)
// SQL 出过错 → 事务必须回滚,不管回调返回什么
if txFailed {
_ = db.Tx.Rollback()
return false
}
if !isSuccess {
err = db.Tx.Rollback()
if err != nil {
+135
View File
@@ -0,0 +1,135 @@
package db
import (
. "code.hoteas.com/golang/hotime/common"
"database/sql"
"testing"
"github.com/sirupsen/logrus"
)
func newTestDB(t *testing.T) *HoTimeDB {
t.Helper()
sqlDB, err := sql.Open("sqlite3", ":memory:")
if err != nil {
t.Fatal(err)
}
logger := logrus.New()
logger.SetLevel(logrus.WarnLevel)
db := &HoTimeDB{
DB: sqlDB,
Type: "sqlite3",
Dialect: &SQLiteDialect{},
LastErr: &Error{Logger: logger},
Log: logger,
}
_, execErr := sqlDB.Exec("CREATE TABLE test_item (id INTEGER PRIMARY KEY, name TEXT, value INTEGER)")
if execErr != nil {
t.Fatal(execErr)
}
_, execErr = sqlDB.Exec("INSERT INTO test_item (id, name, value) VALUES (1, 'foo', 100)")
if execErr != nil {
t.Fatal(execErr)
}
return db
}
func TestAction_TxFailedForceRollback(t *testing.T) {
db := newTestDB(t)
defer db.DB.Close()
result := db.Action(func(tx HoTimeDB) (isSuccess bool) {
tx.Update("test_item", Map{"value": 999}, Map{"id": 1})
tx.Exec("THIS IS INVALID SQL THAT WILL FAIL")
return true
})
if result != false {
t.Fatalf("Action should return false when SQL error occurred, got true")
}
row := db.Get("test_item", "value", Map{"id": 1})
if row == nil {
t.Fatal("failed to read test_item")
}
val := row.GetCeilInt64("value")
if val != 100 {
t.Fatalf("value should be rolled back to 100, got %d", val)
}
}
func TestAction_NormalCommit(t *testing.T) {
db := newTestDB(t)
defer db.DB.Close()
result := db.Action(func(tx HoTimeDB) (isSuccess bool) {
tx.Update("test_item", Map{"value": 200}, Map{"id": 1})
return true
})
if result != true {
t.Fatalf("Action should return true on success, got false")
}
row := db.Get("test_item", "value", Map{"id": 1})
if row == nil {
t.Fatal("failed to read test_item")
}
val := row.GetCeilInt64("value")
if val != 200 {
t.Fatalf("value should be committed to 200, got %d", val)
}
}
func TestAction_NormalRollback(t *testing.T) {
db := newTestDB(t)
defer db.DB.Close()
result := db.Action(func(tx HoTimeDB) (isSuccess bool) {
tx.Update("test_item", Map{"value": 300}, Map{"id": 1})
return false
})
if result != false {
t.Fatalf("Action should return false, got true")
}
row := db.Get("test_item", "value", Map{"id": 1})
if row == nil {
t.Fatal("failed to read test_item")
}
val := row.GetCeilInt64("value")
if val != 100 {
t.Fatalf("value should be rolled back to 100, got %d", val)
}
}
func TestAction_SqlErrorThenReturnTrue_MustRollback(t *testing.T) {
db := newTestDB(t)
defer db.DB.Close()
result := db.Action(func(tx HoTimeDB) (isSuccess bool) {
tx.Update("test_item", Map{"value": 500}, Map{"id": 1})
tx.Exec("INSERT INTO nonexistent_table (x) VALUES (1)")
tx.Update("test_item", Map{"value": 600}, Map{"id": 1})
return true
})
if result != false {
t.Fatalf("Action should return false when SQL error occurred mid-transaction, got true")
}
row := db.Get("test_item", "value", Map{"id": 1})
if row == nil {
t.Fatal("failed to read test_item")
}
val := row.GetCeilInt64("value")
if val != 100 {
t.Fatalf("value should be rolled back to 100 (original), got %d", val)
}
}