From 457510fa09c5117bcd1dc0651253a3c9e56f6c4a Mon Sep 17 00:00:00 2001 From: silverwind Date: Sat, 3 Oct 2026 09:25:47 +0200 Subject: [PATCH] fix: use READ_COMMITTED_SNAPSHOT on MSSQL (#39512) MSSQL's default READ COMMITTED makes reads wait on writers, so the runner pickup deadlocks with concurrent claims, flaking `TestCreateTaskForRunnerConcurrentClaim`. - Enable `READ_COMMITTED_SNAPSHOT` on MSSQL so it reads like PostgreSQL and MySQL - Read the pickup cursor before claiming, a lost claim could skip waiting jobs - Add tests that fail without consistent READ COMMITTED Performance: Writes on MSSQL now also store the previous row version in tempdb, the same versioning cost PostgreSQL and MySQL always pay, and Azure SQL enables it by default. Reads no longer block on writers, and a 32-runner pickup stress test ran 2.5x faster with it. --------- Signed-off-by: wxiaoguang Co-authored-by: wxiaoguang Co-authored-by: Giteabot --- docs/guidelines-backend.md | 4 ++ models/actions/task.go | 11 +++-- models/db/engine_init.go | 14 ++++++ .../actions_concurrent_claim_test.go | 43 +++++++++++++++++++ 4 files changed, 68 insertions(+), 4 deletions(-) diff --git a/docs/guidelines-backend.md b/docs/guidelines-backend.md index 87c67b2771e..66500ac3914 100644 --- a/docs/guidelines-backend.md +++ b/docs/guidelines-backend.md @@ -62,6 +62,10 @@ Operations that must roll back together should run inside `db.WithTx()` (or Functions that participate in a transaction take a `context.Context` as their first parameter so the transaction can be propagated. +PostgreSQL, MySQL and MSSQL (via `READ_COMMITTED_SNAPSHOT`) read the last committed +row version, so reads never wait for writers. Guard read-then-write logic with a +conditional `UPDATE` or a lock. + ### XORM gotchas - Never call `x.Update(exemplar)` without an explicit `WHERE` clause — it updates diff --git a/models/actions/task.go b/models/actions/task.go index 0e922be83bf..838db470f8a 100644 --- a/models/actions/task.go +++ b/models/actions/task.go @@ -298,6 +298,12 @@ func CreateTaskForRunner(ctx context.Context, runner *ActionRunner) (*ActionTask if err := e.Where(cond).Asc("updated", "id").Limit(pickTaskBatchSize).Find(&jobs); err != nil { return nil, false, err } + // A short page means no waiting jobs remain beyond it. + isLastPage := len(jobs) < pickTaskBatchSize + if !isLastPage { + last := jobs[len(jobs)-1] // read before a lost claim bumps Updated + cursorUpdated, cursorID = last.Updated, last.ID + } for _, v := range jobs { if !runner.CanMatchLabels(v.RunsOn) { @@ -313,12 +319,9 @@ func CreateTaskForRunner(ctx context.Context, runner *ActionRunner) (*ActionTask // Another runner claimed this job concurrently; try the next one. } - // A short page means no waiting jobs remain beyond it. - if len(jobs) < pickTaskBatchSize { + if isLastPage { return nil, false, nil } - last := jobs[len(jobs)-1] - cursorUpdated, cursorID = last.Updated, last.ID } } diff --git a/models/db/engine_init.go b/models/db/engine_init.go index 4f46bb33660..043ff6283ed 100644 --- a/models/db/engine_init.go +++ b/models/db/engine_init.go @@ -7,6 +7,7 @@ import ( "context" "database/sql" "fmt" + "time" "gitea.dev/modules/log" "gitea.dev/modules/setting" @@ -109,6 +110,10 @@ func InitEngineWithMigration(ctx context.Context, migrateFunc func(context.Conte preprocessDatabaseCollation(xormEngine) + if setting.Database.Type.IsMSSQL() { + enableMSSQLReadCommittedSnapshot(ctx, xormEngine) + } + // We have to run migrateFunc here in case the user is re-running installation on a previously created DB. // If we do not then table schemas will be changed and there will be conflicts when the migrations run properly. // @@ -131,3 +136,12 @@ func InitEngineWithMigration(ctx context.Context, migrateFunc func(context.Conte return nil } + +// enableMSSQLReadCommittedSnapshot stops MSSQL reads waiting on writers, like PostgreSQL and MySQL +func enableMSSQLReadCommittedSnapshot(ctx context.Context, engine EngineMigration) { + ctx, cancel := context.WithTimeout(ctx, 5*time.Second) // ALTER waits for all other connections to close + defer cancel() + if _, err := engine.Context(ctx).Exec("IF (SELECT is_read_committed_snapshot_on FROM sys.databases WHERE database_id = DB_ID()) = 0 ALTER DATABASE CURRENT SET READ_COMMITTED_SNAPSHOT ON"); err != nil { + log.Error("Unable to set READ_COMMITTED_SNAPSHOT=ON: %v", err) + } +} diff --git a/tests/integration/actions_concurrent_claim_test.go b/tests/integration/actions_concurrent_claim_test.go index c33955352ec..c097dbfd5bf 100644 --- a/tests/integration/actions_concurrent_claim_test.go +++ b/tests/integration/actions_concurrent_claim_test.go @@ -4,16 +4,20 @@ package integration import ( + "context" "sync" "testing" + "time" actions_model "gitea.dev/models/actions" "gitea.dev/models/db" "gitea.dev/models/unittest" + "gitea.dev/modules/setting" "gitea.dev/tests" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" + "xorm.io/builder" ) // minimalWorkflowPayload returns the minimal YAML for a single-job workflow with no steps. @@ -116,3 +120,42 @@ func TestCreateTaskForRunnerConcurrentClaim(t *testing.T) { assert.NotZero(t, updated.TaskID) } } + +func prepareWaitingRunJob(t *testing.T) *actions_model.ActionRunJob { + if setting.Database.Type.IsSQLite3() { + t.Skip("SQLite serializes write transactions") + } + job := &actions_model.ActionRunJob{RepoID: 1, Status: actions_model.StatusWaiting, RunsOn: []string{"ubuntu-latest"}} + require.NoError(t, db.Insert(t.Context(), job)) + return job +} + +func TestCreateTaskForRunnerDuringOpenClaimDoesNotWait(t *testing.T) { + job := prepareWaitingRunJob(t) + + require.NoError(t, db.WithTx(t.Context(), func(ctx context.Context) error { + _, err := db.GetEngine(ctx).ID(job.ID).Cols("name").Update(&actions_model.ActionRunJob{Name: "claiming"}) + require.NoError(t, err) + pickupCtx, cancel := context.WithTimeout(t.Context(), 5*time.Second) + defer cancel() + _, _, err = actions_model.CreateTaskForRunner(pickupCtx, &actions_model.ActionRunner{}) + return err + })) +} + +func TestClaimRunJobAfterConcurrentCancelUpdatesNothing(t *testing.T) { + job := prepareWaitingRunJob(t) + + require.NoError(t, db.WithTx(t.Context(), func(ctx context.Context) error { + claimed, err := actions_model.GetRunJobByRepoAndID(ctx, job.RepoID, job.ID) + require.NoError(t, err) + _, err = db.GetEngine(t.Context()).ID(job.ID).Cols("status").Update(&actions_model.ActionRunJob{Status: actions_model.StatusCancelled}) + require.NoError(t, err) + + claimed.TaskID, claimed.Status = 1, actions_model.StatusRunning + affected, err := actions_model.UpdateRunJob(ctx, claimed, builder.Eq{"task_id": 0, "status": actions_model.StatusWaiting}, "task_id", "status") + require.NoError(t, err) + assert.Zero(t, affected) + return nil + })) +}