mirror of
https://github.com/go-gitea/gitea.git
synced 2026-10-02 15:10:43 +02:00
feat(actions): add build queue view (#38585)
Adds a read-only Actions job queue: running jobs first, then waiting jobs in the order a runner picks them up. It is shown instance-wide in the admin Actions section with owner, repository and status filters, and per repository in the Actions tab. Both lists refresh in place. Pending work is currently only visible per repository and newest-first, so nothing shows what is queued, in which order, or what occupies the runners. Reordering the queue will be proposed separately. A migration adds indexes for the runner pickup query and repository-scoped status lookups. * Fix #34198 <img width="1345" height="451" alt="image" src="https://github.com/user-attachments/assets/7d52ff76-81b4-44e8-b583-d7d89c9dffcd" /> <img width="1809" height="1134" alt="image" src="https://github.com/user-attachments/assets/4d56c0cb-bae7-4ce2-8f3c-75163b2bc7f4" /> --------- Co-authored-by: Zettat123 <zettat123@gmail.com> Co-authored-by: silverwind <me@silverwind.io> Co-authored-by: wxiaoguang <wxiaoguang@gmail.com>
This commit is contained in:
30 files changed
+928
-95
No files matched your search
@@ -30,7 +30,7 @@ import (
|
||||
type ActionRun struct {
|
||||
ID int64
|
||||
Title string
|
||||
RepoID int64 `xorm:"unique(repo_index)"`
|
||||
RepoID int64 `xorm:"unique(repo_index) index(repo_status)"`
|
||||
Repo *repo_model.Repository `xorm:"-"`
|
||||
OwnerID int64 `xorm:"index"`
|
||||
WorkflowID string `xorm:"index"` // the name of workflow file
|
||||
@@ -47,7 +47,7 @@ type ActionRun struct {
|
||||
Event webhook_module.HookEventType // the webhook event that causes the workflow to run
|
||||
EventPayload string `xorm:"LONGTEXT"`
|
||||
TriggerEvent string // the trigger event defined in the `on` configuration of the triggered workflow
|
||||
Status Status `xorm:"index"`
|
||||
Status Status `xorm:"index index(repo_status)"`
|
||||
Version int `xorm:"version default 0"` // Status could be updated concomitantly, so an optimistic lock is needed
|
||||
RawConcurrency string // raw concurrency
|
||||
|
||||
|
||||
@@ -33,7 +33,7 @@ type ActionRunJob struct {
|
||||
ID int64
|
||||
RunID int64 `xorm:"index"`
|
||||
Run *ActionRun `xorm:"-"`
|
||||
RepoID int64 `xorm:"index(repo_concurrency)"`
|
||||
RepoID int64 `xorm:"index(repo_concurrency) index(repo_status)"`
|
||||
Repo *repo_model.Repository `xorm:"-"`
|
||||
OwnerID int64 `xorm:"index"`
|
||||
CommitSHA string `xorm:"index"`
|
||||
@@ -52,10 +52,10 @@ type ActionRunJob struct {
|
||||
Needs []string `xorm:"JSON TEXT"`
|
||||
RunsOn []string `xorm:"JSON TEXT"`
|
||||
|
||||
TaskID int64 // the task created by this job in its own attempt
|
||||
TaskID int64 `xorm:"index(pickup)"` // the task created by this job in its own attempt
|
||||
SourceTaskID int64 `xorm:"NOT NULL DEFAULT 0"` // SourceTaskID points to a historical task when this job reuses an earlier attempt's result.
|
||||
|
||||
Status Status `xorm:"index"`
|
||||
Status Status `xorm:"index index(pickup) index(repo_status)"`
|
||||
|
||||
RawConcurrency string // raw concurrency from job YAML's "concurrency" section
|
||||
|
||||
@@ -131,7 +131,7 @@ type ActionRunJob struct {
|
||||
Started timeutil.TimeStamp
|
||||
Stopped timeutil.TimeStamp
|
||||
Created timeutil.TimeStamp `xorm:"created"`
|
||||
Updated timeutil.TimeStamp `xorm:"updated index"`
|
||||
Updated timeutil.TimeStamp `xorm:"updated index index(pickup)"`
|
||||
}
|
||||
|
||||
// ActionRunAttemptJobIDIndex backs the run-wide AttemptJobID counter, keyed by ActionRun.ID.
|
||||
|
||||
@@ -6,7 +6,12 @@ package actions
|
||||
import (
|
||||
"testing"
|
||||
|
||||
"gitea.dev/models/db"
|
||||
"gitea.dev/models/unittest"
|
||||
"gitea.dev/modules/timeutil"
|
||||
|
||||
"github.com/stretchr/testify/assert"
|
||||
"github.com/stretchr/testify/require"
|
||||
)
|
||||
|
||||
func TestActionJobList_SortMatrixGroupsByName(t *testing.T) {
|
||||
@@ -59,3 +64,49 @@ func TestActionJobList_SortMatrixGroupsByName(t *testing.T) {
|
||||
assert.Equal(t, []string{"only"}, names(jobs))
|
||||
})
|
||||
}
|
||||
|
||||
func TestFindJobQueueJobs(t *testing.T) {
|
||||
require.NoError(t, unittest.PrepareTestDatabase())
|
||||
ctx := t.Context()
|
||||
const repoID int64 = 987654
|
||||
|
||||
insert := func(status Status, taskID int64, reusable bool, updated timeutil.TimeStamp) int64 {
|
||||
job := &ActionRunJob{RepoID: repoID, Status: status, TaskID: taskID, IsReusableCaller: reusable, Updated: updated}
|
||||
_, err := db.GetEngine(ctx).NoAutoTime().Insert(job)
|
||||
require.NoError(t, err)
|
||||
return job.ID
|
||||
}
|
||||
queuedA := insert(StatusWaiting, 0, false, 200)
|
||||
queuedB := insert(StatusWaiting, 0, false, 300)
|
||||
queuedC := insert(StatusWaiting, 0, false, 100)
|
||||
running := insert(StatusRunning, 998, false, 0)
|
||||
cancelling := insert(StatusCancelling, 999, false, 0)
|
||||
insert(StatusWaiting, 999, false, 0)
|
||||
insert(StatusWaiting, 0, true, 0)
|
||||
|
||||
find := func(status Status, page, pageSize int) (ids []int64, total int64) {
|
||||
jobs, total, err := FindJobQueueJobs(ctx, JobQueueOptions{RepoID: repoID, Status: status}, page, pageSize)
|
||||
require.NoError(t, err)
|
||||
for _, job := range jobs {
|
||||
ids = append(ids, job.ID)
|
||||
}
|
||||
return ids, total
|
||||
}
|
||||
|
||||
ids, total := find(StatusUnknown, 1, 10)
|
||||
assert.EqualValues(t, 5, total)
|
||||
assert.Equal(t, []int64{running, cancelling, queuedC, queuedA, queuedB}, ids)
|
||||
|
||||
ids, _ = find(StatusWaiting, 1, 10)
|
||||
assert.Equal(t, []int64{queuedC, queuedA, queuedB}, ids)
|
||||
|
||||
ids, _ = find(StatusRunning, 1, 10)
|
||||
assert.Equal(t, []int64{running, cancelling}, ids)
|
||||
|
||||
ids, _ = find(StatusUnknown, 99, 3)
|
||||
assert.Equal(t, []int64{queuedA, queuedB}, ids)
|
||||
|
||||
repoIDs, err := JobQueueFilterRepoIDs(ctx, 1000)
|
||||
require.NoError(t, err)
|
||||
assert.Contains(t, repoIDs, repoID)
|
||||
}
|
||||
@@ -0,0 +1,81 @@
|
||||
// Copyright 2026 The Gitea Authors. All rights reserved.
|
||||
// SPDX-License-Identifier: MIT
|
||||
|
||||
package actions
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
|
||||
"gitea.dev/models/db"
|
||||
|
||||
"xorm.io/builder"
|
||||
"xorm.io/xorm"
|
||||
)
|
||||
|
||||
// JobQueueOptions scopes the job queue to a repo, an owner or, when both are zero, the instance.
|
||||
type JobQueueOptions struct {
|
||||
RepoID int64
|
||||
OwnerID int64
|
||||
Status Status
|
||||
}
|
||||
|
||||
func (opts JobQueueOptions) session(ctx context.Context) *xorm.Session {
|
||||
// a reusable-workflow caller only tracks its children, it never occupies a runner itself
|
||||
sess := db.GetEngine(ctx).Table("action_run_job").Where(builder.Eq{"`action_run_job`.is_reusable_caller": false})
|
||||
if opts.RepoID > 0 {
|
||||
sess = sess.And(builder.Eq{"`action_run_job`.repo_id": opts.RepoID})
|
||||
}
|
||||
if opts.OwnerID > 0 {
|
||||
sess = sess.Join("INNER", "repository", "repository.id = `action_run_job`.repo_id AND repository.owner_id = ?", opts.OwnerID)
|
||||
}
|
||||
return sess.And(opts.statusCond())
|
||||
}
|
||||
|
||||
var (
|
||||
// keep in sync with CreateTaskForRunner
|
||||
queuedJobsCond = builder.Eq{"`action_run_job`.status": StatusWaiting, "`action_run_job`.task_id": 0}
|
||||
// a cancelling job still occupies its runner
|
||||
runningJobsCond = builder.In("`action_run_job`.status", StatusRunning, StatusCancelling)
|
||||
)
|
||||
|
||||
func (opts JobQueueOptions) statusCond() builder.Cond {
|
||||
switch opts.Status {
|
||||
case StatusRunning:
|
||||
return runningJobsCond
|
||||
case StatusWaiting:
|
||||
return queuedJobsCond
|
||||
default:
|
||||
return builder.Or(runningJobsCond, queuedJobsCond)
|
||||
}
|
||||
}
|
||||
|
||||
// active jobs first by start time, then queued jobs in pickup order
|
||||
var jobQueueOrderBy = fmt.Sprintf(
|
||||
"CASE WHEN `action_run_job`.status IN (%d, %d) THEN 0 ELSE 1 END ASC, CASE WHEN `action_run_job`.status IN (%d, %d) THEN `action_run_job`.started ELSE `action_run_job`.updated END ASC, `action_run_job`.id ASC",
|
||||
StatusRunning, StatusCancelling, StatusRunning, StatusCancelling)
|
||||
|
||||
// FindJobQueueJobs returns one page of the job queue and its total count.
|
||||
func FindJobQueueJobs(ctx context.Context, opts JobQueueOptions, page, pageSize int) ([]*ActionRunJob, int64, error) {
|
||||
total, err := opts.session(ctx).Count(new(ActionRunJob))
|
||||
if err != nil || total == 0 {
|
||||
return nil, total, err
|
||||
}
|
||||
|
||||
// Auto-refresh can shrink the queue under a user still on page 2; show the last page instead of empty.
|
||||
page = min(page, int((total+int64(pageSize)-1)/int64(pageSize)))
|
||||
|
||||
jobs := make([]*ActionRunJob, 0, pageSize)
|
||||
return jobs, total, opts.session(ctx).
|
||||
Cols("`action_run_job`.id", "`action_run_job`.repo_id", "`action_run_job`.name", "`action_run_job`.status", // skip the payload columns
|
||||
"`action_run_job`.run_id", "`action_run_job`.runs_on", "`action_run_job`.updated", "`action_run_job`.started", "`action_run_job`.task_id").
|
||||
OrderBy(jobQueueOrderBy).
|
||||
Limit(pageSize, (page-1)*pageSize).
|
||||
Find(&jobs)
|
||||
}
|
||||
|
||||
// JobQueueFilterRepoIDs returns up to limit ids of the repositories with queued or running jobs.
|
||||
func JobQueueFilterRepoIDs(ctx context.Context, limit int) ([]int64, error) {
|
||||
var ids []int64
|
||||
return ids, JobQueueOptions{}.session(ctx).Distinct("`action_run_job`.repo_id").Limit(limit).Find(&ids)
|
||||
}
|
||||
@@ -160,6 +160,29 @@ func GetTasksMapByIDs(ctx context.Context, ids []int64) (map[int64]*ActionTask,
|
||||
return tasks, db.GetEngine(ctx).In("id", ids).Find(&tasks)
|
||||
}
|
||||
|
||||
// GetTaskRunnerNames returns runner names keyed by task ID without loading task logs.
|
||||
func GetTaskRunnerNames(ctx context.Context, taskIDs []int64) (map[int64]string, error) {
|
||||
names := make(map[int64]string, len(taskIDs))
|
||||
if len(taskIDs) == 0 {
|
||||
return names, nil
|
||||
}
|
||||
var rows []struct {
|
||||
ID int64
|
||||
Name string
|
||||
}
|
||||
err := db.GetEngine(ctx).Table("action_task").
|
||||
Join("INNER", "action_runner", "action_runner.id = action_task.runner_id").
|
||||
In("action_task.id", taskIDs).
|
||||
Select("action_task.id, action_runner.name").Find(&rows)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
for _, row := range rows {
|
||||
names[row.ID] = row.Name
|
||||
}
|
||||
return names, nil
|
||||
}
|
||||
|
||||
func GetRunningTaskByToken(ctx context.Context, token string) (*ActionTask, error) {
|
||||
errNotExist := fmt.Errorf("task with token %q: %w", token, util.ErrNotExist)
|
||||
if token == "" {
|
||||
|
||||
@@ -38,6 +38,18 @@ func TestActionTask_GetRunJobLink(t *testing.T) {
|
||||
assert.Empty(t, (&ActionTask{Job: &ActionRunJob{ID: 42, Run: &ActionRun{ID: 10}}}).GetRunJobLink())
|
||||
}
|
||||
|
||||
func TestGetTaskRunnerNames(t *testing.T) {
|
||||
require.NoError(t, unittest.PrepareTestDatabase())
|
||||
ctx := t.Context()
|
||||
runner := &ActionRunner{Name: "queue-runner"}
|
||||
require.NoError(t, db.Insert(ctx, runner))
|
||||
task := &ActionTask{RunnerID: runner.ID, TokenHash: "queue-test-task"}
|
||||
require.NoError(t, db.Insert(ctx, task))
|
||||
names, err := GetTaskRunnerNames(ctx, []int64{task.ID, 987654321})
|
||||
require.NoError(t, err)
|
||||
assert.Equal(t, map[int64]string{task.ID: runner.Name}, names)
|
||||
}
|
||||
|
||||
func TestMakeTaskStepDisplayName(t *testing.T) {
|
||||
tests := []struct {
|
||||
name string
|
||||
|
||||
Reference in new issue
Block a user