diff --git a/models/actions/run.go b/models/actions/run.go index be291715012..f6c1e005373 100644 --- a/models/actions/run.go +++ b/models/actions/run.go @@ -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) } diff --git a/models/actions/run_attempt.go b/models/actions/run_attempt.go index f592be34cfd..384ecf0ae7f 100644 --- a/models/actions/run_attempt.go +++ b/models/actions/run_attempt.go @@ -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() diff --git a/models/actions/run_job.go b/models/actions/run_job.go index 98247accd69..23c5333b1d9 100644 --- a/models/actions/run_job.go +++ b/models/actions/run_job.go @@ -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 diff --git a/models/actions/run_job_list.go b/models/actions/run_job_list.go index 62d16334540..9ae32335f25 100644 --- a/models/actions/run_job_list.go +++ b/models/actions/run_job_list.go @@ -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)) } diff --git a/services/actions/clear_tasks.go b/services/actions/clear_tasks.go index c7e2e8c0ec0..a2fe4bc1b76 100644 --- a/services/actions/clear_tasks.go +++ b/services/actions/clear_tasks.go @@ -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 } diff --git a/services/actions/clear_tasks_test.go b/services/actions/clear_tasks_test.go index ab654f8985e..923218a9137 100644 --- a/services/actions/clear_tasks_test.go +++ b/services/actions/clear_tasks_test.go @@ -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()) diff --git a/services/actions/job_emitter.go b/services/actions/job_emitter.go index b9e53044319..0b974b10de8 100644 --- a/services/actions/job_emitter.go +++ b/services/actions/job_emitter.go @@ -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 diff --git a/services/actions/job_emitter_test.go b/services/actions/job_emitter_test.go index 0fc4255cc40..6c96abab105 100644 --- a/services/actions/job_emitter_test.go +++ b/services/actions/job_emitter_test.go @@ -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) { diff --git a/services/actions/max_parallel_test.go b/services/actions/max_parallel_test.go index 88bab1f8c61..a39e7ce709e 100644 --- a/services/actions/max_parallel_test.go +++ b/services/actions/max_parallel_test.go @@ -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) { diff --git a/services/actions/reusable_workflow.go b/services/actions/reusable_workflow.go index f2a6f56e67a..1ac89705a56 100644 --- a/services/actions/reusable_workflow.go +++ b/services/actions/reusable_workflow.go @@ -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) } diff --git a/services/actions/run.go b/services/actions/run.go index 6762fda027b..aa33a7f65ce 100644 --- a/services/actions/run.go +++ b/services/actions/run.go @@ -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) } diff --git a/tests/integration/actions_concurrency_test.go b/tests/integration/actions_concurrency_test.go index 50f93959ac3..02f46b2802f 100644 --- a/tests/integration/actions_concurrency_test.go +++ b/tests/integration/actions_concurrency_test.go @@ -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)) diff --git a/tests/integration/actions_reusable_workflow_test.go b/tests/integration/actions_reusable_workflow_test.go index 77c6a0895ce..3f4ccf0880f 100644 --- a/tests/integration/actions_reusable_workflow_test.go +++ b/tests/integration/actions_reusable_workflow_test.go @@ -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) }) }) }