perf(cluster): reduce agent command queue contention

This commit is contained in:
Awuqing
2026-08-09 02:31:44 +08:00
parent ae9623c078
commit b25099b05e
11 changed files with 118 additions and 17 deletions
+9 -1
View File
@@ -4,6 +4,7 @@ import (
"fmt"
"os"
"path/filepath"
"strings"
"backupx/server/internal/config"
"backupx/server/internal/model"
@@ -18,7 +19,14 @@ func Open(cfg config.DatabaseConfig, logger *zap.Logger) (*gorm.DB, error) {
return nil, fmt.Errorf("create database dir: %w", err)
}
db, err := gorm.Open(sqlite.Open(cfg.Path), &gorm.Config{Logger: gormlogger.Default.LogMode(gormlogger.Silent)})
separator := "?"
if strings.Contains(cfg.Path, "?") {
separator = "&"
}
// busy_timeout 减少 Agent 轮询、心跳和任务写入同时发生时的瞬时锁错误。
// 维持默认回滚日志模式,保证当前嵌入式 SQLite 依赖的数据完整性。
dsn := cfg.Path + separator + "_pragma=busy_timeout(5000)"
db, err := gorm.Open(sqlite.Open(dsn), &gorm.Config{Logger: gormlogger.Default.LogMode(gormlogger.Silent)})
if err != nil {
return nil, fmt.Errorf("open sqlite: %w", err)
}
+40
View File
@@ -0,0 +1,40 @@
package database
import (
"path/filepath"
"testing"
"backupx/server/internal/config"
"backupx/server/internal/logger"
)
func TestOpenConfiguresSQLiteForSingleMasterConcurrency(t *testing.T) {
log, err := logger.New(config.LogConfig{Level: "error"})
if err != nil {
t.Fatal(err)
}
db, err := Open(config.DatabaseConfig{Path: filepath.Join(t.TempDir(), "backupx.db")}, log)
if err != nil {
t.Fatal(err)
}
sqlDB, err := db.DB()
if err != nil {
t.Fatal(err)
}
t.Cleanup(func() { _ = sqlDB.Close() })
var journalMode string
if err := db.Raw("PRAGMA journal_mode").Scan(&journalMode).Error; err != nil {
t.Fatal(err)
}
if journalMode != "delete" {
t.Fatalf("journal_mode = %q, want delete", journalMode)
}
var busyTimeout int
if err := db.Raw("PRAGMA busy_timeout").Scan(&busyTimeout).Error; err != nil {
t.Fatal(err)
}
if busyTimeout != 5000 {
t.Fatalf("busy_timeout = %d, want 5000", busyTimeout)
}
}
+16 -16
View File
@@ -4,11 +4,11 @@ import "time"
// AgentCommand 状态常量
const (
AgentCommandStatusPending = "pending" // 待 Agent 拉取
AgentCommandStatusPending = "pending" // 待 Agent 拉取
AgentCommandStatusDispatched = "dispatched" // Agent 已领取,正在执行
AgentCommandStatusSucceeded = "succeeded" // 执行成功
AgentCommandStatusFailed = "failed" // 执行失败
AgentCommandStatusTimeout = "timeout" // 超时未完成
AgentCommandStatusSucceeded = "succeeded" // 执行成功
AgentCommandStatusFailed = "failed" // 执行失败
AgentCommandStatusTimeout = "timeout" // 超时未完成
)
// AgentCommand 类型常量
@@ -36,20 +36,20 @@ const (
)
// AgentCommand 代表 Master 发给某个 Agent 节点的待执行命令。
// 使用简单的数据库队列实现:Agent 通过 token 长轮询拉取本节点 pending 命令,
// 使用简单的数据库队列实现:Agent 通过 token 定期轮询本节点 pending 命令,
// 执行后回写状态与结果。Master 侧通过定时检查把超时的命令标记为 timeout。
type AgentCommand struct {
ID uint `gorm:"primaryKey" json:"id"`
NodeID uint `gorm:"column:node_id;index;not null" json:"nodeId"`
Type string `gorm:"size:32;index;not null" json:"type"`
Status string `gorm:"size:20;index;not null;default:'pending'" json:"status"`
Payload string `gorm:"type:text" json:"payload"` // JSON
Result string `gorm:"type:text" json:"result"` // JSON(成功结果)
ErrorMessage string `gorm:"column:error_message;type:text" json:"errorMessage"`
DispatchedAt *time.Time `gorm:"column:dispatched_at" json:"dispatchedAt,omitempty"`
CompletedAt *time.Time `gorm:"column:completed_at" json:"completedAt,omitempty"`
CreatedAt time.Time `json:"createdAt"`
UpdatedAt time.Time `json:"updatedAt"`
ID uint `gorm:"primaryKey" json:"id"`
NodeID uint `gorm:"column:node_id;not null;index:idx_agent_commands_node_status,priority:1" json:"nodeId"`
Type string `gorm:"size:32;index;not null" json:"type"`
Status string `gorm:"size:20;not null;default:'pending';index:idx_agent_commands_node_status,priority:2;index:idx_agent_commands_status_dispatched,priority:1;index:idx_agent_commands_status_created,priority:1" json:"status"`
Payload string `gorm:"type:text" json:"payload"` // JSON
Result string `gorm:"type:text" json:"result"` // JSON(成功结果)
ErrorMessage string `gorm:"column:error_message;type:text" json:"errorMessage"`
DispatchedAt *time.Time `gorm:"column:dispatched_at;index:idx_agent_commands_status_dispatched,priority:2" json:"dispatchedAt,omitempty"`
CompletedAt *time.Time `gorm:"column:completed_at" json:"completedAt,omitempty"`
CreatedAt time.Time `gorm:"index:idx_agent_commands_status_created,priority:2" json:"createdAt"`
UpdatedAt time.Time `json:"updatedAt"`
}
func (AgentCommand) TableName() string {
@@ -17,12 +17,30 @@ func newTestDB(t *testing.T) *gorm.DB {
if err != nil {
t.Fatalf("open: %v", err)
}
sqlDB, err := db.DB()
if err != nil {
t.Fatalf("get sql database: %v", err)
}
t.Cleanup(func() { _ = sqlDB.Close() })
if err := db.AutoMigrate(&model.AgentCommand{}); err != nil {
t.Fatalf("migrate: %v", err)
}
return db
}
func TestAgentCommandQueueIndexes(t *testing.T) {
db := newTestDB(t)
for _, name := range []string{
"idx_agent_commands_node_status",
"idx_agent_commands_status_dispatched",
"idx_agent_commands_status_created",
} {
if !db.Migrator().HasIndex(&model.AgentCommand{}, name) {
t.Fatalf("missing Agent command queue index %s", name)
}
}
}
func TestAgentCommandRepository_ClaimPending(t *testing.T) {
db := newTestDB(t)
repo := NewAgentCommandRepository(db)
@@ -22,6 +22,11 @@ func newBackupRecordTestRepository(t *testing.T) *GormBackupRecordRepository {
if err != nil {
t.Fatalf("database.Open returned error: %v", err)
}
sqlDB, err := db.DB()
if err != nil {
t.Fatalf("db.DB returned error: %v", err)
}
t.Cleanup(func() { _ = sqlDB.Close() })
storageTarget := &model.StorageTarget{Name: "local", Type: "local_disk", Enabled: true, ConfigCiphertext: "{}", ConfigVersion: 1, LastTestStatus: "unknown"}
if err := db.Create(storageTarget).Error; err != nil {
t.Fatalf("seed storage target error: %v", err)
@@ -21,6 +21,11 @@ func newBackupTaskTestRepository(t *testing.T) *GormBackupTaskRepository {
if err != nil {
t.Fatalf("database.Open returned error: %v", err)
}
sqlDB, err := db.DB()
if err != nil {
t.Fatalf("db.DB returned error: %v", err)
}
t.Cleanup(func() { _ = sqlDB.Close() })
if err := db.Create(&model.StorageTarget{Name: "local", Type: "local_disk", Enabled: true, ConfigCiphertext: "{}", ConfigVersion: 1, LastTestStatus: "unknown"}).Error; err != nil {
t.Fatalf("seed storage target error: %v", err)
}
@@ -19,6 +19,11 @@ func openTestNodeDB(t *testing.T) *gorm.DB {
if err != nil {
t.Fatalf("open sqlite: %v", err)
}
sqlDB, err := db.DB()
if err != nil {
t.Fatalf("get sql database: %v", err)
}
t.Cleanup(func() { _ = sqlDB.Close() })
if err := db.AutoMigrate(&model.Node{}); err != nil {
t.Fatalf("migrate: %v", err)
}
@@ -21,6 +21,11 @@ func newNotificationTestRepository(t *testing.T) *GormNotificationRepository {
if err != nil {
t.Fatalf("database.Open returned error: %v", err)
}
sqlDB, err := db.DB()
if err != nil {
t.Fatalf("db.DB returned error: %v", err)
}
t.Cleanup(func() { _ = sqlDB.Close() })
return NewNotificationRepository(db)
}
@@ -22,6 +22,11 @@ func newOAuthSessionTestRepository(t *testing.T) *GormOAuthSessionRepository {
if err != nil {
t.Fatalf("database.Open returned error: %v", err)
}
sqlDB, err := db.DB()
if err != nil {
t.Fatalf("db.DB returned error: %v", err)
}
t.Cleanup(func() { _ = sqlDB.Close() })
return NewOAuthSessionRepository(db)
}
@@ -22,6 +22,11 @@ func newRestoreRecordTestRepository(t *testing.T) (*GormRestoreRecordRepository,
if err != nil {
t.Fatalf("database.Open returned error: %v", err)
}
sqlDB, err := db.DB()
if err != nil {
t.Fatalf("db.DB returned error: %v", err)
}
t.Cleanup(func() { _ = sqlDB.Close() })
storageTarget := &model.StorageTarget{Name: "local", Type: "local_disk", Enabled: true, ConfigCiphertext: "{}", ConfigVersion: 1, LastTestStatus: "unknown"}
if err := db.Create(storageTarget).Error; err != nil {
t.Fatalf("seed storage target error: %v", err)
@@ -26,6 +26,11 @@ func newStorageTestRepository(t *testing.T) *GormStorageTargetRepository {
if err != nil {
t.Fatalf("database.Open returned error: %v", err)
}
sqlDB, err := db.DB()
if err != nil {
t.Fatalf("db.DB returned error: %v", err)
}
t.Cleanup(func() { _ = sqlDB.Close() })
return NewStorageTargetRepository(db)
}