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 + })) +}