mirror of
https://github.com/go-gitea/gitea.git
synced 2026-10-02 19:39:49 +02:00
feat: Add max-parallel implementation inside the Gitea server
# Conflicts: # models/migrations/migrations.go # models/migrations/v1_26/v325.go # Conflicts: # models/actions/run_job.go # models/actions/runner.go # models/migrations/migrations.go # modules/structs/repo_actions.go # routers/api/actions/runner/runner.go # routers/api/v1/admin/runners.go # routers/api/v1/api.go # routers/api/v1/shared/runners.go # services/actions/run.go # templates/swagger/v1_json.tmpl # Conflicts: # models/migrations/migrations.go
This commit is contained in:
1 parent
a1c60ac854
commit
57400c725e
9 files changed
+958
-3
No files matched your search
@@ -0,0 +1,162 @@
|
||||
// Copyright 2026 The Gitea Authors. All rights reserved.
|
||||
// SPDX-License-Identifier: MIT
|
||||
|
||||
package actions
|
||||
|
||||
import (
|
||||
"context"
|
||||
"testing"
|
||||
|
||||
"code.gitea.io/gitea/models/db"
|
||||
"code.gitea.io/gitea/models/unittest"
|
||||
|
||||
"github.com/stretchr/testify/assert"
|
||||
)
|
||||
|
||||
func TestActionRunJob_MaxParallel(t *testing.T) {
|
||||
assert.NoError(t, unittest.PrepareTestDatabase())
|
||||
ctx := context.Background()
|
||||
|
||||
t.Run("NoMaxParallel", func(t *testing.T) {
|
||||
job := &ActionRunJob{
|
||||
RunID: 1,
|
||||
RepoID: 1,
|
||||
OwnerID: 1,
|
||||
JobID: "test-job-1",
|
||||
Name: "Test Job",
|
||||
Status: StatusWaiting,
|
||||
MaxParallel: 0, // No limit
|
||||
}
|
||||
assert.NoError(t, db.Insert(ctx, job))
|
||||
|
||||
retrieved, err := GetRunJobByID(ctx, job.ID)
|
||||
assert.NoError(t, err)
|
||||
assert.Equal(t, 0, retrieved.MaxParallel)
|
||||
})
|
||||
|
||||
t.Run("WithMaxParallel", func(t *testing.T) {
|
||||
job := &ActionRunJob{
|
||||
RunID: 1,
|
||||
RepoID: 1,
|
||||
OwnerID: 1,
|
||||
JobID: "test-job-2",
|
||||
Name: "Matrix Job",
|
||||
Status: StatusWaiting,
|
||||
MaxParallel: 3,
|
||||
}
|
||||
assert.NoError(t, db.Insert(ctx, job))
|
||||
|
||||
retrieved, err := GetRunJobByID(ctx, job.ID)
|
||||
assert.NoError(t, err)
|
||||
assert.Equal(t, 3, retrieved.MaxParallel)
|
||||
})
|
||||
|
||||
t.Run("MatrixID", func(t *testing.T) {
|
||||
job := &ActionRunJob{
|
||||
RunID: 1,
|
||||
RepoID: 1,
|
||||
OwnerID: 1,
|
||||
JobID: "test-job-3",
|
||||
Name: "Matrix Job with ID",
|
||||
Status: StatusWaiting,
|
||||
MaxParallel: 2,
|
||||
MatrixID: "os:ubuntu,node:16",
|
||||
}
|
||||
assert.NoError(t, db.Insert(ctx, job))
|
||||
|
||||
retrieved, err := GetRunJobByID(ctx, job.ID)
|
||||
assert.NoError(t, err)
|
||||
assert.Equal(t, 2, retrieved.MaxParallel)
|
||||
assert.Equal(t, "os:ubuntu,node:16", retrieved.MatrixID)
|
||||
})
|
||||
|
||||
t.Run("UpdateMaxParallel", func(t *testing.T) {
|
||||
// Create ActionRun first
|
||||
run := &ActionRun{
|
||||
ID: 1,
|
||||
RepoID: 1,
|
||||
OwnerID: 1,
|
||||
Status: StatusRunning,
|
||||
}
|
||||
// Note: This might fail if run already exists from previous tests, but that's okay
|
||||
_ = db.Insert(ctx, run)
|
||||
|
||||
job := &ActionRunJob{
|
||||
RunID: 1,
|
||||
RepoID: 1,
|
||||
OwnerID: 1,
|
||||
JobID: "test-job-4",
|
||||
Name: "Updatable Job",
|
||||
Status: StatusWaiting,
|
||||
MaxParallel: 5,
|
||||
}
|
||||
assert.NoError(t, db.Insert(ctx, job))
|
||||
|
||||
// Update max parallel
|
||||
job.MaxParallel = 10
|
||||
_, err := UpdateRunJob(ctx, job, nil, "max_parallel")
|
||||
assert.NoError(t, err)
|
||||
|
||||
retrieved, err := GetRunJobByID(ctx, job.ID)
|
||||
assert.NoError(t, err)
|
||||
assert.Equal(t, 10, retrieved.MaxParallel)
|
||||
})
|
||||
}
|
||||
|
||||
func TestActionRunJob_MaxParallelEnforcement(t *testing.T) {
|
||||
assert.NoError(t, unittest.PrepareTestDatabase())
|
||||
ctx := context.Background()
|
||||
|
||||
t.Run("EnforceMaxParallel", func(t *testing.T) {
|
||||
runID := int64(5000)
|
||||
jobID := "parallel-enforced-job"
|
||||
maxParallel := 2
|
||||
|
||||
// Create ActionRun first
|
||||
run := &ActionRun{
|
||||
ID: runID,
|
||||
RepoID: 1,
|
||||
OwnerID: 1,
|
||||
Index: 5000,
|
||||
Status: StatusRunning,
|
||||
}
|
||||
assert.NoError(t, db.Insert(ctx, run))
|
||||
|
||||
// Create jobs simulating matrix execution
|
||||
jobs := []*ActionRunJob{
|
||||
{RunID: runID, RepoID: 1, OwnerID: 1, JobID: jobID, Name: "Job 1", Status: StatusRunning, MaxParallel: maxParallel, MatrixID: "version:1"},
|
||||
{RunID: runID, RepoID: 1, OwnerID: 1, JobID: jobID, Name: "Job 2", Status: StatusRunning, MaxParallel: maxParallel, MatrixID: "version:2"},
|
||||
{RunID: runID, RepoID: 1, OwnerID: 1, JobID: jobID, Name: "Job 3", Status: StatusWaiting, MaxParallel: maxParallel, MatrixID: "version:3"},
|
||||
{RunID: runID, RepoID: 1, OwnerID: 1, JobID: jobID, Name: "Job 4", Status: StatusWaiting, MaxParallel: maxParallel, MatrixID: "version:4"},
|
||||
}
|
||||
|
||||
for _, job := range jobs {
|
||||
assert.NoError(t, db.Insert(ctx, job))
|
||||
}
|
||||
|
||||
// Verify running count
|
||||
runningCount, err := CountRunningJobsByWorkflowAndRun(ctx, runID, jobID)
|
||||
assert.NoError(t, err)
|
||||
assert.Equal(t, maxParallel, runningCount, "Should have exactly max-parallel jobs running")
|
||||
|
||||
// Simulate job completion
|
||||
jobs[0].Status = StatusSuccess
|
||||
_, err = UpdateRunJob(ctx, jobs[0], nil, "status")
|
||||
assert.NoError(t, err)
|
||||
|
||||
// Now running count should be 1
|
||||
runningCount, err = CountRunningJobsByWorkflowAndRun(ctx, runID, jobID)
|
||||
assert.NoError(t, err)
|
||||
assert.Equal(t, 1, runningCount)
|
||||
|
||||
// Simulate next job starting
|
||||
jobs[2].Status = StatusRunning
|
||||
_, err = UpdateRunJob(ctx, jobs[2], nil, "status")
|
||||
assert.NoError(t, err)
|
||||
|
||||
// Back to max-parallel
|
||||
runningCount, err = CountRunningJobsByWorkflowAndRun(ctx, runID, jobID)
|
||||
assert.NoError(t, err)
|
||||
assert.Equal(t, maxParallel, runningCount)
|
||||
})
|
||||
}
|
||||
@@ -0,0 +1,95 @@
|
||||
// Copyright 2026 The Gitea Authors. All rights reserved.
|
||||
// SPDX-License-Identifier: MIT
|
||||
|
||||
package actions
|
||||
|
||||
import (
|
||||
"context"
|
||||
"testing"
|
||||
|
||||
"code.gitea.io/gitea/models/db"
|
||||
"code.gitea.io/gitea/models/unittest"
|
||||
|
||||
"github.com/stretchr/testify/assert"
|
||||
)
|
||||
|
||||
func TestActionRunner_Capacity(t *testing.T) {
|
||||
assert.NoError(t, unittest.PrepareTestDatabase())
|
||||
ctx := context.Background()
|
||||
|
||||
t.Run("DefaultCapacity", func(t *testing.T) {
|
||||
runner := &ActionRunner{
|
||||
UUID: "test-uuid-1",
|
||||
Name: "test-runner",
|
||||
OwnerID: 0,
|
||||
RepoID: 0,
|
||||
TokenHash: "hash1",
|
||||
Token: "token1",
|
||||
}
|
||||
assert.NoError(t, db.Insert(ctx, runner))
|
||||
|
||||
// Verify in database
|
||||
retrieved, err := GetRunnerByID(ctx, runner.ID)
|
||||
assert.NoError(t, err)
|
||||
assert.Equal(t, 0, retrieved.Capacity, "Default capacity should be 0 (unlimited)")
|
||||
})
|
||||
|
||||
t.Run("CustomCapacity", func(t *testing.T) {
|
||||
runner := &ActionRunner{
|
||||
UUID: "test-uuid-2",
|
||||
Name: "test-runner-2",
|
||||
OwnerID: 0,
|
||||
RepoID: 0,
|
||||
Capacity: 5,
|
||||
TokenHash: "hash2",
|
||||
Token: "token2",
|
||||
}
|
||||
assert.NoError(t, db.Insert(ctx, runner))
|
||||
|
||||
assert.Equal(t, 5, runner.Capacity)
|
||||
|
||||
// Verify in database
|
||||
retrieved, err := GetRunnerByID(ctx, runner.ID)
|
||||
assert.NoError(t, err)
|
||||
assert.Equal(t, 5, retrieved.Capacity)
|
||||
})
|
||||
|
||||
t.Run("UpdateCapacity", func(t *testing.T) {
|
||||
runner := &ActionRunner{
|
||||
UUID: "test-uuid-3",
|
||||
Name: "test-runner-3",
|
||||
OwnerID: 0,
|
||||
RepoID: 0,
|
||||
Capacity: 1,
|
||||
TokenHash: "hash3",
|
||||
Token: "token3",
|
||||
}
|
||||
assert.NoError(t, db.Insert(ctx, runner))
|
||||
|
||||
// Update capacity
|
||||
runner.Capacity = 10
|
||||
assert.NoError(t, UpdateRunner(ctx, runner, "capacity"))
|
||||
|
||||
// Verify update
|
||||
retrieved, err := GetRunnerByID(ctx, runner.ID)
|
||||
assert.NoError(t, err)
|
||||
assert.Equal(t, 10, retrieved.Capacity)
|
||||
})
|
||||
|
||||
t.Run("ZeroCapacity", func(t *testing.T) {
|
||||
runner := &ActionRunner{
|
||||
UUID: "test-uuid-4",
|
||||
Name: "test-runner-4",
|
||||
OwnerID: 0,
|
||||
RepoID: 0,
|
||||
Capacity: 0, // Unlimited
|
||||
}
|
||||
assert.NoError(t, db.Insert(ctx, runner))
|
||||
|
||||
assert.Equal(t, 0, runner.Capacity)
|
||||
|
||||
retrieved, err := GetRunnerByID(ctx, runner.ID)
|
||||
assert.NoError(t, err)
|
||||
assert.Equal(t, 0, retrieved.Capacity)
|
||||
})
|
||||
}
|
||||
+39
-3
@@ -260,10 +260,26 @@ func CreateTaskForRunner(ctx context.Context, runner *ActionRunner) (*ActionTask
|
||||
var job *ActionRunJob
|
||||
log.Trace("runner labels: %v", runner.AgentLabels)
|
||||
for _, v := range jobs {
|
||||
if runner.CanMatchLabels(v.RunsOn) {
|
||||
job = v
|
||||
break
|
||||
if !runner.CanMatchLabels(v.RunsOn) {
|
||||
continue
|
||||
}
|
||||
|
||||
// Check max-parallel constraint for matrix jobs
|
||||
if v.MaxParallel > 0 {
|
||||
runningCount, err := CountRunningJobsByWorkflowAndRun(ctx, v.RunID, v.JobID)
|
||||
if err != nil {
|
||||
log.Error("Failed to count running jobs for max-parallel check: %v", err)
|
||||
continue
|
||||
}
|
||||
if runningCount >= v.MaxParallel {
|
||||
log.Debug("Job %s (run %d) skipped: %d/%d jobs already running (max-parallel)",
|
||||
v.JobID, v.RunID, runningCount, v.MaxParallel)
|
||||
continue
|
||||
}
|
||||
}
|
||||
|
||||
job = v
|
||||
break
|
||||
}
|
||||
if job == nil {
|
||||
return nil, false, nil
|
||||
@@ -522,3 +538,23 @@ func getTaskIDFromCache(token string) int64 {
|
||||
}
|
||||
return t
|
||||
}
|
||||
|
||||
// CountRunningTasksByRunner counts the number of running tasks assigned to a specific runner
|
||||
func CountRunningTasksByRunner(ctx context.Context, runnerID int64) (int, error) {
|
||||
count, err := db.GetEngine(ctx).
|
||||
Where("runner_id = ?", runnerID).
|
||||
And("status = ?", StatusRunning).
|
||||
Count(new(ActionTask))
|
||||
return int(count), err
|
||||
}
|
||||
|
||||
// CountRunningJobsByWorkflowAndRun counts running jobs for a specific workflow/run combo
|
||||
// Used to enforce max-parallel limits on matrix jobs
|
||||
func CountRunningJobsByWorkflowAndRun(ctx context.Context, runID int64, jobID string) (int, error) {
|
||||
count, err := db.GetEngine(ctx).
|
||||
Where("run_id = ?", runID).
|
||||
And("job_id = ?", jobID).
|
||||
And("status = ?", StatusRunning).
|
||||
Count(new(ActionRunJob))
|
||||
return int(count), err
|
||||
}
|
||||
@@ -0,0 +1,206 @@
|
||||
// Copyright 2026 The Gitea Authors. All rights reserved.
|
||||
// SPDX-License-Identifier: MIT
|
||||
|
||||
package actions
|
||||
|
||||
import (
|
||||
"context"
|
||||
"testing"
|
||||
|
||||
"code.gitea.io/gitea/models/db"
|
||||
"code.gitea.io/gitea/models/unittest"
|
||||
|
||||
"github.com/stretchr/testify/assert"
|
||||
)
|
||||
|
||||
func TestCountRunningTasksByRunner(t *testing.T) {
|
||||
assert.NoError(t, unittest.PrepareTestDatabase())
|
||||
ctx := context.Background()
|
||||
|
||||
t.Run("NoRunningTasks", func(t *testing.T) {
|
||||
count, err := CountRunningTasksByRunner(ctx, 999999)
|
||||
assert.NoError(t, err)
|
||||
assert.Equal(t, 0, count)
|
||||
})
|
||||
|
||||
t.Run("WithRunningTasks", func(t *testing.T) {
|
||||
// Create a runner
|
||||
runner := &ActionRunner{
|
||||
UUID: "test-runner-tasks",
|
||||
Name: "Test Runner",
|
||||
OwnerID: 0,
|
||||
RepoID: 0,
|
||||
TokenHash: "test_hash_tasks",
|
||||
Token: "test_token_tasks",
|
||||
}
|
||||
assert.NoError(t, db.Insert(ctx, runner))
|
||||
|
||||
// Create running tasks
|
||||
task1 := &ActionTask{
|
||||
JobID: 1,
|
||||
RunnerID: runner.ID,
|
||||
Status: StatusRunning,
|
||||
RepoID: 1,
|
||||
OwnerID: 1,
|
||||
TokenHash: "task1_hash",
|
||||
Token: "task1_token",
|
||||
}
|
||||
assert.NoError(t, db.Insert(ctx, task1))
|
||||
|
||||
task2 := &ActionTask{
|
||||
JobID: 2,
|
||||
RunnerID: runner.ID,
|
||||
Status: StatusRunning,
|
||||
RepoID: 1,
|
||||
OwnerID: 1,
|
||||
TokenHash: "task2_hash",
|
||||
Token: "task2_token",
|
||||
}
|
||||
assert.NoError(t, db.Insert(ctx, task2))
|
||||
|
||||
// Count should be 2
|
||||
count, err := CountRunningTasksByRunner(ctx, runner.ID)
|
||||
assert.NoError(t, err)
|
||||
assert.Equal(t, 2, count)
|
||||
})
|
||||
|
||||
t.Run("MixedStatusTasks", func(t *testing.T) {
|
||||
runner := &ActionRunner{
|
||||
UUID: "test-runner-mixed",
|
||||
Name: "Mixed Status Runner",
|
||||
Capacity: 5,
|
||||
TokenHash: "mixed_runner_hash",
|
||||
Token: "mixed_runner_token",
|
||||
}
|
||||
assert.NoError(t, db.Insert(ctx, runner))
|
||||
|
||||
// Create tasks with different statuses
|
||||
statuses := []Status{StatusRunning, StatusSuccess, StatusRunning, StatusFailure, StatusWaiting}
|
||||
for i, status := range statuses {
|
||||
task := &ActionTask{
|
||||
JobID: int64(100 + i),
|
||||
RunnerID: runner.ID,
|
||||
Status: status,
|
||||
RepoID: 1,
|
||||
OwnerID: 1,
|
||||
TokenHash: "mixed_task_hash_" + string(rune('a'+i)),
|
||||
Token: "mixed_task_token_" + string(rune('a'+i)),
|
||||
}
|
||||
assert.NoError(t, db.Insert(ctx, task))
|
||||
}
|
||||
|
||||
// Only 2 running tasks
|
||||
count, err := CountRunningTasksByRunner(ctx, runner.ID)
|
||||
assert.NoError(t, err)
|
||||
assert.Equal(t, 2, count)
|
||||
})
|
||||
}
|
||||
|
||||
func TestCountRunningJobsByWorkflowAndRun(t *testing.T) {
|
||||
assert.NoError(t, unittest.PrepareTestDatabase())
|
||||
ctx := context.Background()
|
||||
|
||||
t.Run("NoRunningJobs", func(t *testing.T) {
|
||||
count, err := CountRunningJobsByWorkflowAndRun(ctx, 999999, "nonexistent")
|
||||
assert.NoError(t, err)
|
||||
assert.Equal(t, 0, count)
|
||||
})
|
||||
|
||||
t.Run("WithRunningJobs", func(t *testing.T) {
|
||||
runID := int64(1000)
|
||||
jobID := "test-job"
|
||||
|
||||
// Create ActionRun first
|
||||
run := &ActionRun{
|
||||
ID: runID,
|
||||
RepoID: 1,
|
||||
OwnerID: 1,
|
||||
Index: 1000,
|
||||
Status: StatusRunning,
|
||||
}
|
||||
assert.NoError(t, db.Insert(ctx, run))
|
||||
|
||||
// Create running jobs
|
||||
for range 3 {
|
||||
job := &ActionRunJob{
|
||||
RunID: runID,
|
||||
RepoID: 1,
|
||||
OwnerID: 1,
|
||||
JobID: jobID,
|
||||
Name: "Test Job",
|
||||
Status: StatusRunning,
|
||||
}
|
||||
assert.NoError(t, db.Insert(ctx, job))
|
||||
}
|
||||
|
||||
count, err := CountRunningJobsByWorkflowAndRun(ctx, runID, jobID)
|
||||
assert.NoError(t, err)
|
||||
assert.Equal(t, 3, count)
|
||||
})
|
||||
|
||||
t.Run("DifferentJobIDs", func(t *testing.T) {
|
||||
runID := int64(2000)
|
||||
|
||||
// Create ActionRun first
|
||||
run := &ActionRun{
|
||||
ID: runID,
|
||||
RepoID: 1,
|
||||
OwnerID: 1,
|
||||
Index: 2000,
|
||||
Status: StatusRunning,
|
||||
}
|
||||
assert.NoError(t, db.Insert(ctx, run))
|
||||
|
||||
// Create jobs with different job IDs
|
||||
for i := range 5 {
|
||||
job := &ActionRunJob{
|
||||
RunID: runID,
|
||||
RepoID: 1,
|
||||
OwnerID: 1,
|
||||
JobID: "job-" + string(rune('A'+i)),
|
||||
Name: "Test Job",
|
||||
Status: StatusRunning,
|
||||
}
|
||||
assert.NoError(t, db.Insert(ctx, job))
|
||||
}
|
||||
|
||||
// Count for specific job ID should be 1
|
||||
count, err := CountRunningJobsByWorkflowAndRun(ctx, runID, "job-A")
|
||||
assert.NoError(t, err)
|
||||
assert.Equal(t, 1, count)
|
||||
})
|
||||
|
||||
t.Run("MatrixJobsWithMaxParallel", func(t *testing.T) {
|
||||
runID := int64(3000)
|
||||
jobID := "matrix-job"
|
||||
maxParallel := 2
|
||||
|
||||
// Create ActionRun first
|
||||
run := &ActionRun{
|
||||
ID: runID,
|
||||
RepoID: 1,
|
||||
OwnerID: 1,
|
||||
Index: 3000,
|
||||
Status: StatusRunning,
|
||||
}
|
||||
assert.NoError(t, db.Insert(ctx, run))
|
||||
|
||||
// Create matrix jobs
|
||||
jobs := []*ActionRunJob{
|
||||
{RunID: runID, RepoID: 1, OwnerID: 1, JobID: jobID, Name: "Job 1", Status: StatusRunning, MaxParallel: maxParallel},
|
||||
{RunID: runID, RepoID: 1, OwnerID: 1, JobID: jobID, Name: "Job 2", Status: StatusRunning, MaxParallel: maxParallel},
|
||||
{RunID: runID, RepoID: 1, OwnerID: 1, JobID: jobID, Name: "Job 3", Status: StatusWaiting, MaxParallel: maxParallel},
|
||||
{RunID: runID, RepoID: 1, OwnerID: 1, JobID: jobID, Name: "Job 4", Status: StatusWaiting, MaxParallel: maxParallel},
|
||||
}
|
||||
|
||||
for _, job := range jobs {
|
||||
assert.NoError(t, db.Insert(ctx, job))
|
||||
}
|
||||
|
||||
// Count running jobs - should be 2 (matching max-parallel)
|
||||
count, err := CountRunningJobsByWorkflowAndRun(ctx, runID, jobID)
|
||||
assert.NoError(t, err)
|
||||
assert.Equal(t, 2, count)
|
||||
assert.Equal(t, maxParallel, count, "Running jobs should equal max-parallel")
|
||||
})
|
||||
}
|
||||
@@ -405,6 +405,7 @@ func prepareMigrationTasks() []*migration {
|
||||
newMigration(328, "Add TokenPermissions column to ActionRunJob", v1_26.AddTokenPermissionsToActionRunJob),
|
||||
newMigration(329, "Add unique constraint for user badge", v1_26.AddUniqueIndexForUserBadge),
|
||||
newMigration(330, "Add name column to webhook", v1_26.AddNameToWebhook),
|
||||
newMigration(331, "Add runner capacity and job max-parallel support", v1_26.AddRunnerCapacityAndJobMaxParallel),
|
||||
}
|
||||
return preparedMigrations
|
||||
}
|
||||
|
||||
Reference in new issue
Block a user