mirror of
https://github.com/go-gitea/gitea.git
synced 2026-10-01 20:59:45 +09:00
fix(actions): preserve admitted jobs and runs in their concurrency group (#39461)
This commit is contained in:
+30
-18
@@ -366,25 +366,42 @@ func UpdateRun(ctx context.Context, run *ActionRun, cols ...string) error {
|
||||
|
||||
type ActionRunIndex db.ResourceIndex
|
||||
|
||||
// GetConcurrentRunAttemptsAndJobs returns run attempts and jobs in the same concurrency group by statuses.
|
||||
func GetConcurrentRunAttemptsAndJobs(ctx context.Context, repoID int64, concurrencyGroup string, status []Status) ([]*ActionRunAttempt, []*ActionRunJob, error) {
|
||||
attempts, err := FindConcurrentRunAttempts(ctx, repoID, concurrencyGroup, status)
|
||||
if err != nil {
|
||||
// expandedCallerCond matches a reusable caller that passed its gate, its status aggregated from its children can still be blocked
|
||||
var expandedCallerCond = builder.Eq{"is_reusable_caller": true, "is_expanded": true}
|
||||
|
||||
// GetConcurrencyHolders returns the run attempts and jobs that passed the gate of the concurrency group and are not done yet.
|
||||
func GetConcurrencyHolders(ctx context.Context, repoID int64, concurrencyGroup string) ([]*ActionRunAttempt, []*ActionRunJob, error) {
|
||||
holding := builder.In("status", StatusWaiting, StatusRunning, StatusCancelling)
|
||||
return findConcurrencyGroupEntries(ctx, repoID, concurrencyGroup, holding, holding.Or(expandedCallerCond.And(builder.Eq{"status": StatusBlocked})))
|
||||
}
|
||||
|
||||
// GetConcurrencyWaiters returns the run attempts and jobs blocked at the gate of the concurrency group.
|
||||
func GetConcurrencyWaiters(ctx context.Context, repoID int64, concurrencyGroup string) ([]*ActionRunAttempt, []*ActionRunJob, error) {
|
||||
blocked := builder.Eq{"status": StatusBlocked}
|
||||
return findConcurrencyGroupEntries(ctx, repoID, concurrencyGroup, blocked, blocked.And(builder.Not{expandedCallerCond}))
|
||||
}
|
||||
|
||||
func findConcurrencyGroupEntries(ctx context.Context, repoID int64, concurrencyGroup string, attemptCond, jobCond builder.Cond) ([]*ActionRunAttempt, []*ActionRunJob, error) {
|
||||
groupCond := builder.Eq{"repo_id": repoID, "concurrency_group": concurrencyGroup}
|
||||
attempts := make([]*ActionRunAttempt, 0)
|
||||
if err := db.GetEngine(ctx).Where(groupCond.And(attemptCond)).Find(&attempts); err != nil {
|
||||
return nil, nil, fmt.Errorf("find run attempts: %w", err)
|
||||
}
|
||||
|
||||
jobs, err := db.Find[ActionRunJob](ctx, &FindRunJobOptions{
|
||||
RepoID: repoID,
|
||||
ConcurrencyGroup: concurrencyGroup,
|
||||
Statuses: status,
|
||||
})
|
||||
if err != nil {
|
||||
jobs := make([]*ActionRunJob, 0)
|
||||
if err := db.GetEngine(ctx).Where(groupCond.And(jobCond)).Find(&jobs); err != nil {
|
||||
return nil, nil, fmt.Errorf("find jobs: %w", err)
|
||||
}
|
||||
|
||||
return attempts, jobs, nil
|
||||
}
|
||||
|
||||
func getConcurrencyEntriesToReplace(ctx context.Context, repoID int64, concurrencyGroup string, cancelInProgress bool) ([]*ActionRunAttempt, []*ActionRunJob, error) {
|
||||
if !cancelInProgress {
|
||||
return GetConcurrencyWaiters(ctx, repoID, concurrencyGroup)
|
||||
}
|
||||
unfinished := builder.In("status", StatusBlocked, StatusWaiting, StatusRunning, StatusCancelling)
|
||||
return findConcurrencyGroupEntries(ctx, repoID, concurrencyGroup, unfinished, unfinished)
|
||||
}
|
||||
|
||||
func CancelPreviousJobsByRunConcurrency(ctx context.Context, attempt *ActionRunAttempt) ([]*ActionRunJob, error) {
|
||||
if attempt.ConcurrencyGroup == "" {
|
||||
return nil, nil
|
||||
@@ -392,12 +409,7 @@ func CancelPreviousJobsByRunConcurrency(ctx context.Context, attempt *ActionRunA
|
||||
|
||||
var jobsToCancel []*ActionRunJob
|
||||
|
||||
statusFindOption := []Status{StatusWaiting, StatusBlocked}
|
||||
if attempt.ConcurrencyCancel {
|
||||
statusFindOption = append(statusFindOption, StatusRunning)
|
||||
statusFindOption = append(statusFindOption, StatusCancelling)
|
||||
}
|
||||
attempts, jobs, err := GetConcurrentRunAttemptsAndJobs(ctx, attempt.RepoID, attempt.ConcurrencyGroup, statusFindOption)
|
||||
attempts, jobs, err := getConcurrencyEntriesToReplace(ctx, attempt.RepoID, attempt.ConcurrencyGroup, attempt.ConcurrencyCancel)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("find concurrent runs and jobs: %w", err)
|
||||
}
|
||||
|
||||
@@ -147,17 +147,6 @@ func findPassThroughAttemptIDs(ctx context.Context, attemptIDs []int64) ([]int64
|
||||
Find(&passThroughAttemptIDs)
|
||||
}
|
||||
|
||||
// FindConcurrentRunAttempts returns attempts in the given concurrency group and status set.
|
||||
// Results are unordered; callers must not depend on any particular row order.
|
||||
func FindConcurrentRunAttempts(ctx context.Context, repoID int64, concurrencyGroup string, statuses []Status) ([]*ActionRunAttempt, error) {
|
||||
attempts := make([]*ActionRunAttempt, 0)
|
||||
sess := db.GetEngine(ctx).Where("repo_id=? AND concurrency_group=?", repoID, concurrencyGroup)
|
||||
if len(statuses) > 0 {
|
||||
sess = sess.In("status", statuses)
|
||||
}
|
||||
return attempts, sess.Find(&attempts)
|
||||
}
|
||||
|
||||
func UpdateRunAttempt(ctx context.Context, attempt *ActionRunAttempt, cols ...string) error {
|
||||
if slices.Contains(cols, "status") && attempt.Started.IsZero() && attempt.Status.IsRunning() {
|
||||
attempt.Started = timeutil.TimeStampNow()
|
||||
|
||||
@@ -724,6 +724,20 @@ func CancelPreviousJobs(ctx context.Context, repoID int64, ref, workflowID strin
|
||||
return cancelledJobs, nil
|
||||
}
|
||||
|
||||
// GetAncestorCallerIDs returns the IDs of the reusable workflow callers the job is nested in.
|
||||
func GetAncestorCallerIDs(ctx context.Context, job *ActionRunJob) (container.Set[int64], error) {
|
||||
ids := make(container.Set[int64])
|
||||
for parentID := job.ParentJobID; parentID != 0; {
|
||||
parent, err := GetRunJobByRunAndID(ctx, job.RunID, parentID)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("load caller %d: %w", parentID, err)
|
||||
}
|
||||
ids.Add(parent.ID)
|
||||
parentID = parent.ParentJobID
|
||||
}
|
||||
return ids, nil
|
||||
}
|
||||
|
||||
func CancelPreviousJobsByJobConcurrency(ctx context.Context, job *ActionRunJob) (jobsToCancel []*ActionRunJob, _ error) {
|
||||
if job.RawConcurrency == "" {
|
||||
return nil, nil
|
||||
@@ -735,16 +749,15 @@ func CancelPreviousJobsByJobConcurrency(ctx context.Context, job *ActionRunJob)
|
||||
return nil, nil
|
||||
}
|
||||
|
||||
statusFindOption := []Status{StatusWaiting, StatusBlocked}
|
||||
if job.ConcurrencyCancel {
|
||||
statusFindOption = append(statusFindOption, StatusRunning)
|
||||
statusFindOption = append(statusFindOption, StatusCancelling)
|
||||
}
|
||||
attempts, jobs, err := GetConcurrentRunAttemptsAndJobs(ctx, job.RepoID, job.ConcurrencyGroup, statusFindOption)
|
||||
attempts, jobs, err := getConcurrencyEntriesToReplace(ctx, job.RepoID, job.ConcurrencyGroup, job.ConcurrencyCancel)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("find concurrent runs and jobs: %w", err)
|
||||
}
|
||||
jobs = slices.DeleteFunc(jobs, func(j *ActionRunJob) bool { return j.ID == job.ID })
|
||||
callerIDs, err := GetAncestorCallerIDs(ctx, job)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
jobs = slices.DeleteFunc(jobs, func(j *ActionRunJob) bool { return j.ID == job.ID || callerIDs.Contains(j.ID) })
|
||||
jobsToCancel = append(jobsToCancel, jobs...)
|
||||
|
||||
// cancel runs in the same concurrency group
|
||||
|
||||
@@ -91,15 +91,14 @@ func (jobs ActionJobList) LoadAttributes(ctx context.Context, withRepo bool) err
|
||||
|
||||
type FindRunJobOptions struct {
|
||||
db.ListOptions
|
||||
RunID int64
|
||||
RunAttemptID optional.Option[int64] // use optional to allow filtering by zero (legacy jobs have run_attempt_id=0)
|
||||
RepoID int64
|
||||
OwnerID int64
|
||||
CommitSHA string
|
||||
Statuses []Status
|
||||
UpdatedBefore timeutil.TimeStamp
|
||||
ConcurrencyGroup string
|
||||
OrderBy db.SearchOrderBy
|
||||
RunID int64
|
||||
RunAttemptID optional.Option[int64] // use optional to allow filtering by zero (legacy jobs have run_attempt_id=0)
|
||||
RepoID int64
|
||||
OwnerID int64
|
||||
CommitSHA string
|
||||
Statuses []Status
|
||||
UpdatedBefore timeutil.TimeStamp
|
||||
OrderBy db.SearchOrderBy
|
||||
// AccessibleRepoIDsSubQuery, when non-nil, restricts results to the repo IDs selected by the
|
||||
// subquery (the caller's accessible repos). A nil value means no restriction. Using a subquery
|
||||
// instead of a materialized ID slice avoids exceeding DB parameter limits for large owners.
|
||||
@@ -131,12 +130,6 @@ func (opts FindRunJobOptions) ToConds() builder.Cond {
|
||||
if opts.UpdatedBefore > 0 {
|
||||
cond = cond.And(builder.Lt{"`action_run_job`.updated": opts.UpdatedBefore})
|
||||
}
|
||||
if opts.ConcurrencyGroup != "" {
|
||||
if opts.RepoID == 0 {
|
||||
panic("Invalid FindRunJobOptions: repo_id is required")
|
||||
}
|
||||
cond = cond.And(builder.Eq{"`action_run_job`.concurrency_group": opts.ConcurrencyGroup})
|
||||
}
|
||||
if opts.AccessibleRepoIDsSubQuery != nil {
|
||||
cond = cond.And(builder.In("`action_run_job`.repo_id", opts.AccessibleRepoIDsSubQuery))
|
||||
}
|
||||
|
||||
@@ -6,6 +6,7 @@ package actions
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"slices"
|
||||
"time"
|
||||
|
||||
actions_model "gitea.dev/models/actions"
|
||||
@@ -71,28 +72,33 @@ func shouldBlockJobByConcurrency(ctx context.Context, job *actions_model.ActionR
|
||||
return false, nil
|
||||
}
|
||||
|
||||
attempts, jobs, err := actions_model.GetConcurrentRunAttemptsAndJobs(ctx, job.RepoID, job.ConcurrencyGroup, []actions_model.Status{actions_model.StatusRunning, actions_model.StatusCancelling})
|
||||
attempts, jobs, err := actions_model.GetConcurrencyHolders(ctx, job.RepoID, job.ConcurrencyGroup)
|
||||
if err != nil {
|
||||
return false, fmt.Errorf("GetConcurrentRunAttemptsAndJobs: %w", err)
|
||||
return false, fmt.Errorf("GetConcurrencyHolders: %w", err)
|
||||
}
|
||||
|
||||
return len(attempts) > 0 || len(jobs) > 0, nil
|
||||
callerIDs, err := actions_model.GetAncestorCallerIDs(ctx, job)
|
||||
if err != nil {
|
||||
return false, err
|
||||
}
|
||||
// the job's own attempt and callers may declare the same group, they must not block their own job
|
||||
return slices.ContainsFunc(attempts, func(a *actions_model.ActionRunAttempt) bool { return a.ID != job.RunAttemptID }) ||
|
||||
slices.ContainsFunc(jobs, func(j *actions_model.ActionRunJob) bool { return !callerIDs.Contains(j.ID) }), nil
|
||||
}
|
||||
|
||||
// PrepareToStartJobWithConcurrency prepares a job to start by its evaluated concurrency group and cancelling previous jobs if necessary.
|
||||
// It returns the new status of the job (either StatusBlocked or StatusWaiting), any cancelled jobs, and any error encountered during the process.
|
||||
func PrepareToStartJobWithConcurrency(ctx context.Context, job *actions_model.ActionRunJob) (actions_model.Status, []*actions_model.ActionRunJob, error) {
|
||||
shouldBlock, err := shouldBlockJobByConcurrency(ctx, job)
|
||||
if err != nil {
|
||||
return actions_model.StatusBlocked, nil, err
|
||||
}
|
||||
|
||||
// even if the current job is blocked, we still need to cancel previous "waiting/blocked" jobs in the same concurrency group
|
||||
// cancel before checking, so the jobs this cancellation finishes no longer hold the group
|
||||
jobs, err := actions_model.CancelPreviousJobsByJobConcurrency(ctx, job)
|
||||
if err != nil {
|
||||
return actions_model.StatusBlocked, nil, fmt.Errorf("CancelPreviousJobsByJobConcurrency: %w", err)
|
||||
}
|
||||
|
||||
shouldBlock, err := shouldBlockJobByConcurrency(ctx, job)
|
||||
if err != nil {
|
||||
return actions_model.StatusBlocked, nil, err
|
||||
}
|
||||
|
||||
return util.Iif(shouldBlock, actions_model.StatusBlocked, actions_model.StatusWaiting), jobs, nil
|
||||
}
|
||||
|
||||
@@ -101,28 +107,28 @@ func shouldBlockRunByConcurrency(ctx context.Context, attempt *actions_model.Act
|
||||
return false, nil
|
||||
}
|
||||
|
||||
attempts, jobs, err := actions_model.GetConcurrentRunAttemptsAndJobs(ctx, attempt.RepoID, attempt.ConcurrencyGroup, []actions_model.Status{actions_model.StatusRunning, actions_model.StatusCancelling})
|
||||
attempts, jobs, err := actions_model.GetConcurrencyHolders(ctx, attempt.RepoID, attempt.ConcurrencyGroup)
|
||||
if err != nil {
|
||||
return false, fmt.Errorf("find concurrent runs and jobs: %w", err)
|
||||
}
|
||||
|
||||
return len(attempts) > 0 || len(jobs) > 0, nil
|
||||
// the run's own attempt and jobs may declare the same group, they must not block their own run
|
||||
return slices.ContainsFunc(attempts, func(a *actions_model.ActionRunAttempt) bool { return a.RunID != attempt.RunID }) ||
|
||||
slices.ContainsFunc(jobs, func(j *actions_model.ActionRunJob) bool { return j.RunID != attempt.RunID }), nil
|
||||
}
|
||||
|
||||
// PrepareToStartRunWithConcurrency prepares a run attempt to start by its evaluated concurrency group and cancelling previous jobs if necessary.
|
||||
// It returns the new status of the run attempt (either StatusBlocked or StatusWaiting), any cancelled jobs, and any error encountered during the process.
|
||||
func PrepareToStartRunWithConcurrency(ctx context.Context, attempt *actions_model.ActionRunAttempt) (actions_model.Status, []*actions_model.ActionRunJob, error) {
|
||||
shouldBlock, err := shouldBlockRunByConcurrency(ctx, attempt)
|
||||
if err != nil {
|
||||
return actions_model.StatusBlocked, nil, err
|
||||
}
|
||||
|
||||
// even if the current run is blocked, we still need to cancel previous "waiting/blocked" jobs in the same concurrency group
|
||||
jobs, err := actions_model.CancelPreviousJobsByRunConcurrency(ctx, attempt)
|
||||
if err != nil {
|
||||
return actions_model.StatusBlocked, nil, fmt.Errorf("CancelPreviousJobsByRunConcurrency: %w", err)
|
||||
}
|
||||
|
||||
shouldBlock, err := shouldBlockRunByConcurrency(ctx, attempt)
|
||||
if err != nil {
|
||||
return actions_model.StatusBlocked, nil, err
|
||||
}
|
||||
|
||||
return util.Iif(shouldBlock, actions_model.StatusBlocked, actions_model.StatusWaiting), jobs, nil
|
||||
}
|
||||
|
||||
|
||||
@@ -19,7 +19,7 @@ import (
|
||||
"github.com/stretchr/testify/require"
|
||||
)
|
||||
|
||||
func createConflictingCancellingJob(t *testing.T, concurrencyGroup string, runIndex int64) *actions_model.ActionRunJob {
|
||||
func createRunAttempt(t *testing.T, runIndex int64, concurrencyGroup string, status actions_model.Status) (*actions_model.ActionRun, *actions_model.ActionRunAttempt) {
|
||||
t.Helper()
|
||||
|
||||
run := &actions_model.ActionRun{
|
||||
@@ -29,7 +29,7 @@ func createConflictingCancellingJob(t *testing.T, concurrencyGroup string, runIn
|
||||
WorkflowID: "test.yml",
|
||||
Index: runIndex,
|
||||
Ref: "refs/heads/main",
|
||||
Status: actions_model.StatusBlocked,
|
||||
Status: status,
|
||||
}
|
||||
require.NoError(t, db.Insert(t.Context(), run))
|
||||
|
||||
@@ -38,11 +38,18 @@ func createConflictingCancellingJob(t *testing.T, concurrencyGroup string, runIn
|
||||
RunID: run.ID,
|
||||
Attempt: 1,
|
||||
TriggerUserID: run.TriggerUserID,
|
||||
Status: actions_model.StatusBlocked,
|
||||
Status: status,
|
||||
ConcurrencyGroup: concurrencyGroup,
|
||||
}
|
||||
require.NoError(t, db.Insert(t.Context(), attempt))
|
||||
|
||||
return run, attempt
|
||||
}
|
||||
|
||||
func createConflictingCancellingJob(t *testing.T, concurrencyGroup string, runIndex int64) *actions_model.ActionRunJob {
|
||||
t.Helper()
|
||||
|
||||
run, attempt := createRunAttempt(t, runIndex, concurrencyGroup, actions_model.StatusBlocked)
|
||||
job := &actions_model.ActionRunJob{
|
||||
RunID: run.ID,
|
||||
RunAttemptID: attempt.ID,
|
||||
@@ -107,6 +114,70 @@ func TestShouldBlockJobByConcurrency_CancellingJobBlocks(t *testing.T) {
|
||||
assert.True(t, shouldBlock)
|
||||
}
|
||||
|
||||
func TestShouldBlockJobByConcurrency_OwnAttemptDoesNotBlock(t *testing.T) {
|
||||
require.NoError(t, unittest.PrepareTestDatabase())
|
||||
|
||||
const concurrencyGroup = "test-own-attempt-does-not-block"
|
||||
_, attempt := createRunAttempt(t, 9906, concurrencyGroup, actions_model.StatusWaiting)
|
||||
job := &actions_model.ActionRunJob{RunAttemptID: attempt.ID, RepoID: attempt.RepoID, ConcurrencyGroup: concurrencyGroup}
|
||||
shouldBlock, err := shouldBlockJobByConcurrency(t.Context(), job)
|
||||
require.NoError(t, err)
|
||||
assert.False(t, shouldBlock)
|
||||
|
||||
createRunAttempt(t, 9907, concurrencyGroup, actions_model.StatusWaiting)
|
||||
shouldBlock, err = shouldBlockJobByConcurrency(t.Context(), job)
|
||||
require.NoError(t, err)
|
||||
assert.True(t, shouldBlock)
|
||||
}
|
||||
|
||||
func TestPrepareToStartJobWithConcurrency_ExpandedCallerHoldsGroup(t *testing.T) {
|
||||
require.NoError(t, unittest.PrepareTestDatabase())
|
||||
|
||||
const concurrencyGroup = "test-expanded-caller-holds-group"
|
||||
_, callerAttempt := createRunAttempt(t, 9908, "", actions_model.StatusBlocked)
|
||||
caller := &actions_model.ActionRunJob{
|
||||
RunID: callerAttempt.RunID, RunAttemptID: callerAttempt.ID, RepoID: callerAttempt.RepoID, Status: actions_model.StatusBlocked,
|
||||
IsReusableCaller: true, IsExpanded: true, ConcurrencyGroup: concurrencyGroup,
|
||||
}
|
||||
require.NoError(t, db.Insert(t.Context(), caller))
|
||||
|
||||
status, cancelled, err := PrepareToStartJobWithConcurrency(t.Context(), &actions_model.ActionRunJob{
|
||||
RepoID: caller.RepoID, RawConcurrency: concurrencyGroup, IsConcurrencyEvaluated: true, ConcurrencyGroup: concurrencyGroup,
|
||||
})
|
||||
require.NoError(t, err)
|
||||
assert.Equal(t, actions_model.StatusBlocked, status)
|
||||
assert.Empty(t, cancelled)
|
||||
}
|
||||
|
||||
func TestPrepareToStartWithConcurrency_WaitingJobHoldsGroup(t *testing.T) {
|
||||
require.NoError(t, unittest.PrepareTestDatabase())
|
||||
|
||||
const concurrencyGroup = "test-waiting-job-holds-group"
|
||||
_, attempt := createRunAttempt(t, 9905, "", actions_model.StatusWaiting)
|
||||
previousJob := &actions_model.ActionRunJob{
|
||||
RunID: attempt.RunID, RunAttemptID: attempt.ID, RepoID: attempt.RepoID, Status: actions_model.StatusWaiting, ConcurrencyGroup: concurrencyGroup,
|
||||
}
|
||||
require.NoError(t, db.Insert(t.Context(), previousJob))
|
||||
|
||||
newJob := &actions_model.ActionRunJob{RepoID: attempt.RepoID, RawConcurrency: concurrencyGroup, IsConcurrencyEvaluated: true, ConcurrencyGroup: concurrencyGroup}
|
||||
status, cancelled, err := PrepareToStartJobWithConcurrency(t.Context(), newJob)
|
||||
require.NoError(t, err)
|
||||
assert.Equal(t, actions_model.StatusBlocked, status)
|
||||
assert.Empty(t, cancelled)
|
||||
|
||||
status, cancelled, err = PrepareToStartRunWithConcurrency(t.Context(), &actions_model.ActionRunAttempt{RepoID: attempt.RepoID, ConcurrencyGroup: concurrencyGroup})
|
||||
require.NoError(t, err)
|
||||
assert.Equal(t, actions_model.StatusBlocked, status)
|
||||
assert.Empty(t, cancelled)
|
||||
|
||||
newJob.ConcurrencyCancel = true
|
||||
status, cancelled, err = PrepareToStartJobWithConcurrency(t.Context(), newJob)
|
||||
require.NoError(t, err)
|
||||
assert.Equal(t, actions_model.StatusWaiting, status)
|
||||
require.Len(t, cancelled, 1)
|
||||
assert.Equal(t, previousJob.ID, cancelled[0].ID)
|
||||
}
|
||||
|
||||
func TestShouldBlockRunByConcurrency_CancellingJobBlocks(t *testing.T) {
|
||||
assert.NoError(t, unittest.PrepareTestDatabase())
|
||||
|
||||
|
||||
@@ -157,27 +157,23 @@ func findConcurrencyWaiterToWake(ctx context.Context, repoID, excludeRunID int64
|
||||
return 0, nil
|
||||
}
|
||||
|
||||
// The slot should be free before any waiter can proceed.
|
||||
holderAttempts, holderJobs, err := actions_model.GetConcurrentRunAttemptsAndJobs(ctx, repoID, concurrencyGroup, []actions_model.Status{actions_model.StatusRunning, actions_model.StatusCancelling})
|
||||
if err != nil {
|
||||
return 0, fmt.Errorf("find concurrency-group holders: %w", err)
|
||||
}
|
||||
if len(holderAttempts) > 0 || len(holderJobs) > 0 {
|
||||
return 0, nil
|
||||
}
|
||||
|
||||
cAttempts, cJobs, err := actions_model.GetConcurrentRunAttemptsAndJobs(ctx, repoID, concurrencyGroup, []actions_model.Status{actions_model.StatusBlocked})
|
||||
cAttempts, cJobs, err := actions_model.GetConcurrencyWaiters(ctx, repoID, concurrencyGroup)
|
||||
if err != nil {
|
||||
return 0, fmt.Errorf("find blocked concurrent runs: %w", err)
|
||||
}
|
||||
// the waiter's own gate decides, it ignores holders that are the waiter's own run or callers
|
||||
for _, a := range cAttempts {
|
||||
if a.RunID != excludeRunID {
|
||||
return a.RunID, nil
|
||||
if blocked, err := shouldBlockRunByConcurrency(ctx, a); err != nil || !blocked {
|
||||
return a.RunID, err
|
||||
}
|
||||
}
|
||||
}
|
||||
for _, j := range cJobs {
|
||||
if j.RunID != excludeRunID {
|
||||
return j.RunID, nil
|
||||
if blocked, err := shouldBlockJobByConcurrency(ctx, j); err != nil || !blocked {
|
||||
return j.RunID, err
|
||||
}
|
||||
}
|
||||
}
|
||||
return 0, nil
|
||||
@@ -366,9 +362,9 @@ func checkJobsOfCurrentRunAttempt(ctx context.Context, run *actions_model.Action
|
||||
}
|
||||
|
||||
result.UpdatedJobs = append(result.UpdatedJobs, resolver.matrixUpdatedJobs...)
|
||||
// Caller and matrix expansion both insert Pending or Blocked jobs, which only a follow-up pass resolves.
|
||||
// Caller and matrix expansion insert Pending or Blocked jobs and a deferred gate leaves a job Blocked, only a follow-up pass resolves them.
|
||||
// Like the caller's children, matrix siblings are left out of result.Jobs and picked up there.
|
||||
if expandedAnyCaller || resolver.matrixChanged {
|
||||
if expandedAnyCaller || resolver.matrixChanged || resolver.gateDeferred {
|
||||
result.RunIDsToReEmit = append(result.RunIDsToReEmit, run.ID)
|
||||
}
|
||||
result.CancelledJobs = append(result.CancelledJobs, resolver.cancelledJobs...)
|
||||
@@ -413,6 +409,9 @@ type jobStatusResolver struct {
|
||||
// matrixUpdatedJobs holds jobs whose status matrix expansion persisted itself, so they are
|
||||
// notified like the ones the caller updates from the resolved status map.
|
||||
matrixUpdatedJobs []*actions_model.ActionRunJob
|
||||
// admittedGroups are the groups a job was admitted to in this pass, which the gate only sees in the database once the pass commits
|
||||
admittedGroups []string
|
||||
gateDeferred bool
|
||||
}
|
||||
|
||||
func newJobStatusResolver(jobs actions_model.ActionJobList, vars map[string]string) *jobStatusResolver {
|
||||
@@ -602,12 +601,22 @@ func (r *jobStatusResolver) resolve(ctx context.Context) (map[int64]actions_mode
|
||||
log.Debug("updateConcurrencyEvaluationForJobWithNeeds failed, this job will stay blocked: job: %d, err: %v", id, err)
|
||||
continue
|
||||
}
|
||||
if actionRunJob.ConcurrencyGroup != "" && slices.Contains(r.admittedGroups, actionRunJob.ConcurrencyGroup) {
|
||||
r.gateDeferred = true
|
||||
continue
|
||||
}
|
||||
|
||||
newStatus, cancelledJobs, err := PrepareToStartJobWithConcurrency(ctx, actionRunJob)
|
||||
if err != nil {
|
||||
log.Error("ShouldBlockJobByConcurrency failed, this job will stay blocked: job: %d, err: %v", id, err)
|
||||
} else {
|
||||
r.cancelledJobs = append(r.cancelledJobs, cancelledJobs...)
|
||||
for _, cancelled := range cancelledJobs {
|
||||
if sibling, ok := r.jobMap[cancelled.ID]; ok { // a sibling replaced here must not reach the gate again in this pass
|
||||
sibling.Status = cancelled.Status
|
||||
r.statuses[cancelled.ID] = cancelled.Status
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
if newStatus == actions_model.StatusWaiting && !slots.take(actionRunJob) {
|
||||
@@ -616,6 +625,7 @@ func (r *jobStatusResolver) resolve(ctx context.Context) (map[int64]actions_mode
|
||||
|
||||
if newStatus != actions_model.StatusBlocked {
|
||||
ret[id] = newStatus
|
||||
r.admittedGroups = append(r.admittedGroups, actionRunJob.ConcurrencyGroup)
|
||||
}
|
||||
}
|
||||
return ret, nil
|
||||
|
||||
@@ -576,6 +576,33 @@ jobs:
|
||||
assert.Equal(t, actions_model.StatusBlocked, refreshed.Status)
|
||||
}
|
||||
|
||||
func Test_checkJobsOfCurrentRunAttempt_SameGroupSiblingsGateInTurn(t *testing.T) {
|
||||
require.NoError(t, unittest.PrepareTestDatabase())
|
||||
ctx := t.Context()
|
||||
|
||||
run, attempt := createRunAttempt(t, 9915, "", actions_model.StatusBlocked)
|
||||
run.LatestAttemptID = attempt.ID
|
||||
jobs := make([]*actions_model.ActionRunJob, 3)
|
||||
for i, jobID := range []string{"a", "b", "c"} {
|
||||
jobs[i] = &actions_model.ActionRunJob{
|
||||
RunID: run.ID, RunAttemptID: attempt.ID, AttemptJobID: int64(i + 1), RepoID: run.RepoID, OwnerID: run.OwnerID,
|
||||
JobID: jobID, Name: jobID, Status: actions_model.StatusBlocked, RawConcurrency: "group: siblings\n", WorkflowPayload: minimalWorkflowPayload(jobID),
|
||||
}
|
||||
require.NoError(t, db.Insert(ctx, jobs[i]))
|
||||
}
|
||||
|
||||
result, err := checkJobsOfCurrentRunAttempt(ctx, run)
|
||||
require.NoError(t, err)
|
||||
assert.Equal(t, []int64{run.ID}, result.RunIDsToReEmit)
|
||||
|
||||
result, err = checkJobsOfCurrentRunAttempt(ctx, run)
|
||||
require.NoError(t, err)
|
||||
assert.Empty(t, result.RunIDsToReEmit)
|
||||
for i, status := range []actions_model.Status{actions_model.StatusWaiting, actions_model.StatusBlocked, actions_model.StatusCancelled} {
|
||||
assert.Equal(t, status, unittest.AssertExistsAndLoadBean(t, &actions_model.ActionRunJob{ID: jobs[i].ID}).Status)
|
||||
}
|
||||
}
|
||||
|
||||
func Test_checkJobsOfCurrentRunAttempt_NeedApprovalKeepsJobsBlocked(t *testing.T) {
|
||||
assert.NoError(t, unittest.PrepareTestDatabase())
|
||||
ctx := t.Context()
|
||||
@@ -707,12 +734,25 @@ func Test_findConcurrencyWaiterToWake(t *testing.T) {
|
||||
assert.NoError(t, err)
|
||||
assert.Equal(t, int64(0), id)
|
||||
|
||||
// Held group "held-cg" (a running holder) has a blocked waiter, but nothing is woken while held.
|
||||
seed(99704, "held-cg", actions_model.StatusRunning)
|
||||
seed(99704, "held-cg", actions_model.StatusWaiting)
|
||||
seed(99705, "held-cg", actions_model.StatusBlocked)
|
||||
id, err = findConcurrencyWaiterToWake(ctx, repoID, 0, "held-cg")
|
||||
assert.NoError(t, err)
|
||||
assert.Equal(t, int64(0), id)
|
||||
|
||||
own := seed(99706, "", actions_model.StatusBlocked)
|
||||
caller := &actions_model.ActionRunJob{RunID: own.ID, RepoID: repoID, Status: actions_model.StatusBlocked, IsReusableCaller: true, IsExpanded: true, ConcurrencyGroup: "own-cg"}
|
||||
assert.NoError(t, db.Insert(ctx, caller))
|
||||
assert.NoError(t, db.Insert(ctx, &actions_model.ActionRunJob{RunID: own.ID, RepoID: repoID, ParentJobID: caller.ID, Status: actions_model.StatusBlocked, ConcurrencyGroup: "own-cg"}))
|
||||
id, err = findConcurrencyWaiterToWake(ctx, repoID, 0, "own-cg")
|
||||
assert.NoError(t, err)
|
||||
assert.Equal(t, own.ID, id)
|
||||
|
||||
ownRun := seed(99707, "own-run-cg", actions_model.StatusBlocked)
|
||||
assert.NoError(t, db.Insert(ctx, &actions_model.ActionRunJob{RunID: ownRun.ID, RepoID: repoID, Status: actions_model.StatusBlocked, IsReusableCaller: true, IsExpanded: true, ConcurrencyGroup: "own-run-cg"}))
|
||||
id, err = findConcurrencyWaiterToWake(ctx, repoID, 0, "own-run-cg")
|
||||
assert.NoError(t, err)
|
||||
assert.Equal(t, ownRun.ID, id)
|
||||
}
|
||||
|
||||
func Test_maxParallelReusableCallerLifecycle(t *testing.T) {
|
||||
|
||||
@@ -226,6 +226,7 @@ func TestPrepareRunAndInsert_MaxParallelStarvedSkipsConcurrency(t *testing.T) {
|
||||
|
||||
holder = unittest.AssertExistsAndLoadBean(t, &actions_model.ActionRunJob{ID: holder.ID})
|
||||
assert.Equal(t, actions_model.StatusRunning, holder.Status, "the starved job must not cancel the group holder")
|
||||
assert.Empty(t, unittest.AssertExistsAndLoadBean(t, &actions_model.ActionRunJob{RunID: run.ID, Status: actions_model.StatusBlocked}).ConcurrencyGroup)
|
||||
}
|
||||
|
||||
func Test_jobStatusResolver_MaxParallelStarvedSkipsConcurrency(t *testing.T) {
|
||||
|
||||
@@ -28,6 +28,7 @@ import (
|
||||
"gitea.dev/modules/util"
|
||||
"gitea.dev/services/convert"
|
||||
|
||||
"go.yaml.in/yaml/v4"
|
||||
"xorm.io/builder"
|
||||
)
|
||||
|
||||
@@ -424,6 +425,13 @@ func insertCallerChildren(ctx context.Context, run *actions_model.ActionRun, att
|
||||
child.IsReusableCaller = true
|
||||
child.CallUses = parsedChild.Uses
|
||||
}
|
||||
if parsedChild.RawConcurrency != nil {
|
||||
rawConcurrency, err := yaml.Marshal(parsedChild.RawConcurrency)
|
||||
if err != nil {
|
||||
return fmt.Errorf("marshal raw concurrency of child %q under caller %d: %w", jobID, caller.ID, err)
|
||||
}
|
||||
child.RawConcurrency = string(rawConcurrency)
|
||||
}
|
||||
if err := db.Insert(ctx, child); err != nil {
|
||||
return fmt.Errorf("insert child %q under caller %d: %w", jobID, caller.ID, err)
|
||||
}
|
||||
|
||||
@@ -252,8 +252,8 @@ func insertRunJob(ctx context.Context, run *actions_model.ActionRun, runAttempt
|
||||
}
|
||||
runJob.RawConcurrency = string(rawConcurrency)
|
||||
|
||||
// the job emitter evaluates it for jobs with `needs`, a skipped job never takes part
|
||||
if len(needs) == 0 && runJob.Status != actions_model.StatusSkipped {
|
||||
// a job enters its group at its gate, ApproveRuns gates an approval-blocked one with this evaluation
|
||||
if runJob.Status == actions_model.StatusWaiting && slots.available(runJob) || len(needs) == 0 && run.NeedApproval {
|
||||
if err := EvaluateJobConcurrencyFillModel(ctx, run, runAttempt, runJob, vars, inputs); err != nil {
|
||||
return nil, nil, false, fmt.Errorf("evaluate job concurrency: %w", err)
|
||||
}
|
||||
|
||||
@@ -977,14 +977,13 @@ jobs:
|
||||
req = NewRequest(t, "POST", fmt.Sprintf("/%s/%s/actions/runs/%d/rerun", user2.Name, apiRepo.Name, run3.ID))
|
||||
_ = session.MakeRequest(t, req, http.StatusOK)
|
||||
|
||||
assert.Equal(t, actions_model.StatusBlocked, unittest.AssertExistsAndLoadBean(t, &actions_model.ActionRun{ID: run3.ID}).Status)
|
||||
runner.execTask(t, runner.fetchTask(t), &mockTaskOutcome{result: runnerv1.Result_RESULT_SUCCESS})
|
||||
task6 := runner.fetchTask(t)
|
||||
_, _, run3_2 := getTaskAndJobAndRunByTaskID(t, task6.Id)
|
||||
assert.Equal(t, run3.ID, run3_2.ID)
|
||||
assert.Equal(t, actions_model.StatusRunning, run3_2.Status)
|
||||
assert.Equal(t, "workflow-dispatch-v1.22", getRunConcurrencyGroup(t, run3))
|
||||
|
||||
run2_2 = unittest.AssertExistsAndLoadBean(t, &actions_model.ActionRun{ID: run2_2.ID})
|
||||
assert.Equal(t, actions_model.StatusCancelled, run2_2.Status) // cancelled by run3
|
||||
})
|
||||
}
|
||||
|
||||
@@ -1118,12 +1117,11 @@ jobs:
|
||||
req = NewRequest(t, "POST", fmt.Sprintf("/%s/%s/actions/runs/%d/jobs/%d/rerun", user2.Name, apiRepo.Name, run3.ID, job3.ID))
|
||||
_ = session.MakeRequest(t, req, http.StatusOK)
|
||||
|
||||
assert.Equal(t, actions_model.StatusBlocked, unittest.AssertExistsAndLoadBean(t, &actions_model.ActionRun{ID: run3.ID}).Status)
|
||||
runner.execTask(t, runner.fetchTask(t), &mockTaskOutcome{result: runnerv1.Result_RESULT_SUCCESS})
|
||||
task6 := runner.fetchTask(t)
|
||||
_, _, run3 = getTaskAndJobAndRunByTaskID(t, task6.Id)
|
||||
assert.Equal(t, "workflow-dispatch-v1.22", getRunConcurrencyGroup(t, run3))
|
||||
|
||||
run2_2 = unittest.AssertExistsAndLoadBean(t, &actions_model.ActionRun{ID: run2_2.ID})
|
||||
assert.Equal(t, actions_model.StatusCancelled, run2_2.Status) // cancelled by run3
|
||||
})
|
||||
}
|
||||
|
||||
@@ -1360,9 +1358,8 @@ jobs:
|
||||
w3Run := unittest.AssertExistsAndLoadBean(t, &actions_model.ActionRun{RepoID: repo.ID, WorkflowID: "concurrent-workflow-3.yml"})
|
||||
w3j1Job := unittest.AssertExistsAndLoadBean(t, &actions_model.ActionRunJob{RunID: w3Run.ID, JobID: "wf3-job1"})
|
||||
assert.Equal(t, actions_model.StatusBlocked, w3j1Job.Status)
|
||||
// wf2-job1 is cancelled by wf3-job1
|
||||
w2j1Job = unittest.AssertExistsAndLoadBean(t, &actions_model.ActionRunJob{ID: w2j1Job.ID})
|
||||
assert.Equal(t, actions_model.StatusCancelled, w2j1Job.Status)
|
||||
assert.Equal(t, actions_model.StatusBlocked, w2j1Job.Status)
|
||||
|
||||
// exec wf1-job1
|
||||
runner1.execTask(t, w1j1Task, &mockTaskOutcome{
|
||||
@@ -1403,6 +1400,8 @@ jobs:
|
||||
|
||||
// fetch wf4-job1
|
||||
w4j1Task := runner2.fetchTask(t)
|
||||
_, w2j1Job, _ = getTaskAndJobAndRunByTaskID(t, runner1.fetchTask(t).Id)
|
||||
assert.Equal(t, "wf2-job1", w2j1Job.JobID)
|
||||
// all tasks have been fetched
|
||||
runner1.fetchNoTask(t)
|
||||
runner2.fetchNoTask(t)
|
||||
@@ -1410,7 +1409,7 @@ jobs:
|
||||
_, w2j2Job, w2Run = getTaskAndJobAndRunByTaskID(t, w2j2Task.Id)
|
||||
// wf2-job2 is cancelled because wf4-job1's cancel-in-progress is true
|
||||
assert.Equal(t, actions_model.StatusCancelled, w2j2Job.Status)
|
||||
assert.Equal(t, actions_model.StatusCancelled, w2Run.Status)
|
||||
assert.Equal(t, actions_model.StatusRunning, w2Run.Status)
|
||||
_, w4j1Job, w4Run := getTaskAndJobAndRunByTaskID(t, w4j1Task.Id)
|
||||
assert.Equal(t, "job-group-2", w4j1Job.ConcurrencyGroup)
|
||||
assert.Equal(t, "workflow-group-2", getRunConcurrencyGroup(t, w4Run))
|
||||
|
||||
@@ -911,7 +911,7 @@ jobs:
|
||||
assert.Equal(t, actions_model.StatusSuccess, run.Status)
|
||||
})
|
||||
|
||||
t.Run("No-needs caller if evaluated inline: false skips, true expands", func(t *testing.T) {
|
||||
t.Run("No-needs callers: if evaluated inline, called job concurrency serializes them", func(t *testing.T) {
|
||||
// A no-needs reusable-workflow caller is processed inline during InsertRun, where its own
|
||||
// `if:` is now evaluated before expansion:
|
||||
// - a false `if:` skips the caller without inserting any children, and the skip is
|
||||
@@ -928,6 +928,8 @@ on:
|
||||
jobs:
|
||||
inner:
|
||||
runs-on: ubuntu-latest
|
||||
concurrency:
|
||||
group: called-job
|
||||
steps:
|
||||
- run: echo inner
|
||||
`)
|
||||
@@ -946,6 +948,9 @@ jobs:
|
||||
if: ${{ true }}
|
||||
uses: ./.gitea/workflows/lib.yaml
|
||||
|
||||
will_run_too:
|
||||
uses: ./.gitea/workflows/lib.yaml
|
||||
|
||||
after_skip:
|
||||
needs: [will_skip]
|
||||
runs-on: ubuntu-latest
|
||||
@@ -970,11 +975,16 @@ jobs:
|
||||
willRun := unittest.AssertExistsAndLoadBean(t, &actions_model.ActionRunJob{RunID: runID, JobID: "will_run"})
|
||||
assert.True(t, willRun.IsReusableCaller)
|
||||
assert.True(t, willRun.IsExpanded)
|
||||
innerChild := unittest.AssertExistsAndLoadBean(t, &actions_model.ActionRunJob{RunID: runID, JobID: "inner"})
|
||||
assert.Equal(t, willRun.ID, innerChild.ParentJobID)
|
||||
unittest.AssertExistsAndLoadBean(t, &actions_model.ActionRunJob{RunID: runID, JobID: "inner", ParentJobID: willRun.ID})
|
||||
|
||||
afterSkip := unittest.AssertExistsAndLoadBean(t, &actions_model.ActionRunJob{RunID: runID, JobID: "after_skip"})
|
||||
assert.Equal(t, actions_model.StatusSkipped, afterSkip.Status)
|
||||
|
||||
unittest.AssertCount(t, &actions_model.ActionRunJob{RunID: runID, JobID: "inner", Status: actions_model.StatusBlocked}, 1)
|
||||
runner := newMockRunner()
|
||||
runner.registerAsRepoRunner(t, repo.OwnerName, repo.Name, "mock-runner", []string{"ubuntu-latest"}, false)
|
||||
runner.execTask(t, runner.fetchTask(t), &mockTaskOutcome{result: runnerv1.Result_RESULT_SUCCESS})
|
||||
runner.fetchTask(t)
|
||||
})
|
||||
})
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user