enhance(actions): add pending job status and align job statuses with GitHub (#39376)

This commit is contained in:
silverwind
2026-09-27 20:30:42 +02:00
committed by GitHub
parent 0d09986790
commit 9e7b302bd3
45 changed files with 397 additions and 189 deletions
+1 -1
View File
@@ -174,7 +174,7 @@ func stopTasks(ctx context.Context, opts actions_model.FindTaskOptions) error {
// CancelAbandonedJobs cancels jobs that have not been picked by any runner for a long time
func CancelAbandonedJobs(ctx context.Context) error {
abandonedJobs, err := db.Find[actions_model.ActionRunJob](ctx, actions_model.FindRunJobOptions{
Statuses: []actions_model.Status{actions_model.StatusWaiting, actions_model.StatusBlocked},
Statuses: []actions_model.Status{actions_model.StatusWaiting, actions_model.StatusBlocked, actions_model.StatusPending},
UpdatedBefore: timeutil.TimeStampNow().AddDuration(-setting.Actions.AbandonedJobTimeout),
})
if err != nil {
+49 -6
View File
@@ -53,6 +53,10 @@ func CreateCommitStatusForRunJobs(ctx context.Context, run *actions_model.Action
scopedPrefix = actions_model.ScopedStatusContextPrefix(ctx, run.WorkflowRepoID)
}
var pendingFilter *pendingJobFilter
if slices.ContainsFunc(jobs, func(job *actions_model.ActionRunJob) bool { return job.Status.IsPending() && !job.IsMatrixDeferred }) {
pendingFilter = newPendingJobFilter(ctx, run)
}
for _, job := range jobs {
// A deferred-matrix placeholder's name changes when it expands, so a status created while it
// waits would be orphaned. The emitter reloads the jobs after expanding and creates them
@@ -61,7 +65,7 @@ func CreateCommitStatusForRunJobs(ctx context.Context, run *actions_model.Action
if job.IsMatrixDeferred && !job.Status.IsDone() {
continue
}
if err = createCommitStatus(ctx, run.Repo, event, commitID, scopedPrefix, run, job); err != nil {
if err = createCommitStatus(ctx, run.Repo, event, commitID, scopedPrefix, run, job, pendingFilter); err != nil {
log.Error("Failed to create commit status for job %d: %v", job.ID, err)
}
}
@@ -151,7 +155,7 @@ func getCommitStatusEventNameAndCommitID(run *actions_model.ActionRun) (event, c
return event, commitID, nil
}
func createCommitStatus(ctx context.Context, repo *repo_model.Repository, event, commitID, scopedPrefix string, run *actions_model.ActionRun, job *actions_model.ActionRunJob) error {
func createCommitStatus(ctx context.Context, repo *repo_model.Repository, event, commitID, scopedPrefix string, run *actions_model.ActionRun, job *actions_model.ActionRunJob, pendingFilter *pendingJobFilter) error {
displayName := actions_module.WorkflowDisplayName(run.WorkflowID, job.WorkflowPayload)
ctxName := actions_module.WorkflowStatusContextName(displayName, job.Name, event) // git_model.NewCommitStatus also trims spaces
if run.IsScopedRun {
@@ -160,7 +164,35 @@ func createCommitStatus(ctx context.Context, repo *repo_model.Repository, event,
ctxName = actions_module.ScopedWorkflowStatusContextName(scopedPrefix, displayName, job.Name, event)
}
targetURL := fmt.Sprintf("%s/jobs/%d", run.Link(), job.ID)
return createWorkflowCommitStatus(ctx, repo, commitID, ctxName, run.WorkflowID, toCommitStatus(job.Status), targetURL, toCommitStatusDescription(job))
return createWorkflowCommitStatus(ctx, repo, commitID, ctxName, run.WorkflowID, toCommitStatus(job.Status), targetURL, toCommitStatusDescription(job), pendingFilter.onlyReplace(job, ctxName))
}
// pendingJobFilter keeps optional Pending jobs from posting new statuses
type pendingJobFilter struct {
requiredGlobs []glob.Glob
}
func newPendingJobFilter(ctx context.Context, run *actions_model.ActionRun) *pendingJobFilter {
rules, err := git_model.FindRepoProtectedBranchRules(ctx, run.RepoID)
if err != nil {
log.Error("FindRepoProtectedBranchRules: %v", err)
return nil
}
if slices.ContainsFunc(rules, func(rule *git_model.ProtectedBranch) bool {
return rule.EnableStatusCheck && len(rule.StatusCheckContexts) == 0
}) {
return nil
}
requiredGlobs, err := requiredStatusContextGlobs(ctx, run.Repo, rules)
if err != nil {
log.Error("requiredStatusContextGlobs: %v", err)
return nil
}
return &pendingJobFilter{requiredGlobs: requiredGlobs}
}
func (f *pendingJobFilter) onlyReplace(job *actions_model.ActionRunJob, ctxName string) bool {
return f != nil && job.Status.IsPending() && !slices.ContainsFunc(f.requiredGlobs, func(gp glob.Glob) bool { return gp.Match(ctxName) })
}
// getAllRequiredStatusContextGlobs returns the compiled globs of every status-check context required in the repo:
@@ -170,6 +202,10 @@ func getAllRequiredStatusContextGlobs(ctx context.Context, repo *repo_model.Repo
if err != nil {
return nil, fmt.Errorf("FindRepoProtectedBranchRules: %w", err)
}
return requiredStatusContextGlobs(ctx, repo, rules)
}
func requiredStatusContextGlobs(ctx context.Context, repo *repo_model.Repository, rules git_model.ProtectedBranchRules) ([]glob.Glob, error) {
required, err := pull_service.EffectiveRequiredContexts(ctx, repo, rules...)
if err != nil {
return nil, fmt.Errorf("EffectiveRequiredContexts: %w", err)
@@ -246,7 +282,7 @@ func CreateSkippedCommitStatusForFilteredWorkflow(ctx context.Context, repo *rep
continue
}
// "Skipped" mirrors toCommitStatusDescription for StatusSkipped.
if err := createWorkflowCommitStatus(ctx, repo, commitID, ctxName, workflowID, commitstatus.CommitStatusSkipped, "", "Skipped"); err != nil {
if err := createWorkflowCommitStatus(ctx, repo, commitID, ctxName, workflowID, commitstatus.CommitStatusSkipped, "", "Skipped", false); err != nil {
return err
}
}
@@ -254,7 +290,7 @@ func CreateSkippedCommitStatusForFilteredWorkflow(ctx context.Context, repo *rep
}
// createWorkflowCommitStatus posts the commit status for one workflow-job context.
func createWorkflowCommitStatus(ctx context.Context, repo *repo_model.Repository, commitID, ctxName, workflowID string, state commitstatus.CommitStatusState, targetURL, description string) error {
func createWorkflowCommitStatus(ctx context.Context, repo *repo_model.Repository, commitID, ctxName, workflowID string, state commitstatus.CommitStatusState, targetURL, description string, onlyReplace bool) error {
// Mix the workflow file path into the hash so two workflow files that
// share the same `name:` and job name produce distinct commit statuses
// even though they render identically — matching GitHub's behavior
@@ -280,14 +316,19 @@ func createWorkflowCommitStatus(ctx context.Context, repo *repo_model.Repository
break
}
}
hasPrevious := false
for _, v := range statuses {
if v.ContextHash == ctxHash {
if v.State == state && v.TargetURL == targetURL && v.Description == description {
return nil
}
hasPrevious = true
break
}
}
if onlyReplace && !hasPrevious {
return nil
}
creator := user_model.NewActionsUser()
status := git_model.CommitStatus{
@@ -322,6 +363,8 @@ func toCommitStatusDescription(job *actions_model.ActionRunJob) string {
return "Waiting to run"
case actions_model.StatusBlocked:
return "Blocked by required conditions"
case actions_model.StatusPending:
return "Waiting for needed jobs"
default:
return fmt.Sprintf("Unknown status: %d", job.Status)
}
@@ -333,7 +376,7 @@ func toCommitStatus(status actions_model.Status) commitstatus.CommitStatusState
return commitstatus.CommitStatusSuccess
case actions_model.StatusFailure, actions_model.StatusCancelled:
return commitstatus.CommitStatusFailure
case actions_model.StatusWaiting, actions_model.StatusBlocked, actions_model.StatusRunning, actions_model.StatusCancelling:
case actions_model.StatusWaiting, actions_model.StatusBlocked, actions_model.StatusPending, actions_model.StatusRunning, actions_model.StatusCancelling:
return commitstatus.CommitStatusPending
case actions_model.StatusSkipped:
return commitstatus.CommitStatusSkipped
+40 -10
View File
@@ -35,6 +35,7 @@ func TestCommitStatusDescription(t *testing.T) {
{actions_model.StatusRunning, 0, 0, "In progress"},
{actions_model.StatusWaiting, 0, 0, "Waiting to run"},
{actions_model.StatusBlocked, 0, 0, "Blocked by required conditions"},
{actions_model.StatusPending, 0, 0, "Waiting for needed jobs"},
{actions_model.StatusUnknown, 0, 0, "Unknown status: 0"},
}
for _, tc := range cases {
@@ -71,7 +72,7 @@ func TestCreateCommitStatus_Dedupe(t *testing.T) {
expectedContext := "status-dedupe-test.yaml / status-dedupe-job (push)"
expectedTargetURL := run.Link() + "/jobs/99002"
require.NoError(t, createCommitStatus(t.Context(), repo, "push", commit.ID.String(), "", run, job))
require.NoError(t, createCommitStatus(t.Context(), repo, "push", commit.ID.String(), "", run, job, nil))
statuses := findCommitStatusesForContext(t, repo.ID, commit.ID.String(), expectedContext)
require.Len(t, statuses, 1)
@@ -80,7 +81,7 @@ func TestCreateCommitStatus_Dedupe(t *testing.T) {
assert.Equal(t, expectedTargetURL, statuses[0].TargetURL)
job.Status = actions_model.StatusRunning
require.NoError(t, createCommitStatus(t.Context(), repo, "push", commit.ID.String(), "", run, job))
require.NoError(t, createCommitStatus(t.Context(), repo, "push", commit.ID.String(), "", run, job, nil))
statuses = findCommitStatusesForContext(t, repo.ID, commit.ID.String(), expectedContext)
require.Len(t, statuses, 2)
@@ -89,17 +90,46 @@ func TestCreateCommitStatus_Dedupe(t *testing.T) {
assert.Equal(t, "In progress", statuses[1].Description)
assert.Equal(t, expectedTargetURL, statuses[1].TargetURL)
require.NoError(t, createCommitStatus(t.Context(), repo, "push", commit.ID.String(), "", run, job))
require.NoError(t, createCommitStatus(t.Context(), repo, "push", commit.ID.String(), "", run, job, nil))
statuses = findCommitStatusesForContext(t, repo.ID, commit.ID.String(), expectedContext)
assert.Len(t, statuses, 2)
job.Status = actions_model.StatusSuccess
require.NoError(t, createCommitStatus(t.Context(), repo, "push", commit.ID.String(), "", run, job))
require.NoError(t, createCommitStatus(t.Context(), repo, "push", commit.ID.String(), "", run, job, nil))
statuses = findCommitStatusesForContext(t, repo.ID, commit.ID.String(), expectedContext)
require.Len(t, statuses, 3)
assert.Equal(t, commitstatus.CommitStatusSuccess, statuses[2].State)
}
func TestCreateCommitStatus_HidesOptionalPendingJobs(t *testing.T) {
require.NoError(t, unittest.PrepareTestDatabase())
repo := unittest.AssertExistsAndLoadBean(t, &repo_model.Repository{ID: 4})
branch := unittest.AssertExistsAndLoadBean(t, &git_model.Branch{RepoID: repo.ID, Name: repo.DefaultBranch})
run := &actions_model.ActionRun{ID: 99101, RepoID: repo.ID, Repo: repo, WorkflowID: "ci.yaml"}
deploy := &actions_model.ActionRunJob{ID: 99102, RunID: run.ID, RepoID: repo.ID, Name: "deploy", Status: actions_model.StatusPending}
postDeploy := func(pending *pendingJobFilter) []*git_model.CommitStatus {
require.NoError(t, createCommitStatus(t.Context(), repo, "push", branch.CommitID, "", run, deploy, pending))
return findCommitStatusesForContext(t, repo.ID, branch.CommitID, "ci.yaml / deploy (push)")
}
pending := newPendingJobFilter(t.Context(), run)
assert.Empty(t, postDeploy(pending))
deploy.Status = actions_model.StatusSuccess
assert.Len(t, postDeploy(nil), 1)
deploy.Status = actions_model.StatusPending
assert.Len(t, postDeploy(pending), 2)
require.NoError(t, db.Insert(t.Context(), &git_model.ProtectedBranch{RepoID: repo.ID, RuleName: "main", EnableStatusCheck: true, StatusCheckContexts: []string{"ci.yaml / deploy*"}}))
pending = newPendingJobFilter(t.Context(), run)
assert.False(t, pending.onlyReplace(deploy, "ci.yaml / deploy (push)"))
assert.True(t, pending.onlyReplace(deploy, "other / deploy (push)"))
require.NoError(t, db.Insert(t.Context(), &git_model.ProtectedBranch{RepoID: repo.ID, RuleName: "release", EnableStatusCheck: true}))
assert.Nil(t, newPendingJobFilter(t.Context(), run))
}
func TestGetCommitActionsStatusMap(t *testing.T) {
assert.NoError(t, unittest.PrepareTestDatabase())
@@ -125,7 +155,7 @@ func TestGetCommitActionsStatusMap(t *testing.T) {
RunID: run.ID, RepoID: repo.ID, OwnerID: repo.OwnerID, Name: tc.jobName, Status: tc.status,
}
require.NoError(t, db.Insert(t.Context(), job))
require.NoError(t, createCommitStatus(t.Context(), repo, "push", branch.CommitID, "", run, job))
require.NoError(t, createCommitStatus(t.Context(), repo, "push", branch.CommitID, "", run, job, nil))
}
statuses, err := git_model.GetLatestCommitStatus(t.Context(), repo.ID, branch.CommitID, db.ListOptionsAll)
@@ -184,7 +214,7 @@ jobs:
WorkflowPayload: payload,
}
require.NoError(t, db.Insert(t.Context(), job))
require.NoError(t, createCommitStatus(t.Context(), repo, "pull_request", branch.CommitID, "", run, job))
require.NoError(t, createCommitStatus(t.Context(), repo, "pull_request", branch.CommitID, "", run, job, nil))
}
statuses, err := git_model.GetLatestCommitStatus(t.Context(), repo.ID, branch.CommitID, db.ListOptionsAll)
@@ -242,7 +272,7 @@ func TestCreateCommitStatus_LegacyHashRecovery(t *testing.T) {
Name: "my-job", Status: actions_model.StatusSuccess,
}
require.NoError(t, db.Insert(t.Context(), job))
require.NoError(t, createCommitStatus(t.Context(), repo, "push", branch.CommitID, "", run, job))
require.NoError(t, createCommitStatus(t.Context(), repo, "push", branch.CommitID, "", run, job, nil))
latest, err := git_model.GetLatestCommitStatus(t.Context(), repo.ID, branch.CommitID, db.ListOptionsAll)
require.NoError(t, err)
@@ -299,7 +329,7 @@ func TestCreateCommitStatus_LegacyHashExternalNotAdopted(t *testing.T) {
Name: "my-job", Status: actions_model.StatusSuccess,
}
require.NoError(t, db.Insert(t.Context(), job))
require.NoError(t, createCommitStatus(t.Context(), repo, "push", branch.CommitID, "", run, job))
require.NoError(t, createCommitStatus(t.Context(), repo, "push", branch.CommitID, "", run, job, nil))
latest, err := git_model.GetLatestCommitStatus(t.Context(), repo.ID, branch.CommitID, db.ListOptionsAll)
require.NoError(t, err)
@@ -351,7 +381,7 @@ func TestCreateCommitStatus_UnnamedWorkflowUsesFileName(t *testing.T) {
`),
}
require.NoError(t, db.Insert(t.Context(), job))
require.NoError(t, createCommitStatus(t.Context(), repo, "push", branch.CommitID, "", run, job))
require.NoError(t, createCommitStatus(t.Context(), repo, "push", branch.CommitID, "", run, job, nil))
statuses := findCommitStatusesForContext(t, repo.ID, branch.CommitID, tc.workflowID+" / my-test (push)")
require.Len(t, statuses, 1)
@@ -406,7 +436,7 @@ jobs:
if run.IsScopedRun {
scopedPrefix = actions_model.ScopedStatusContextPrefix(t.Context(), run.WorkflowRepoID)
}
require.NoError(t, createCommitStatus(t.Context(), consumer, "push", branch.CommitID, scopedPrefix, run, job))
require.NoError(t, createCommitStatus(t.Context(), consumer, "push", branch.CommitID, scopedPrefix, run, job, nil))
}
// repo-level Context is the bare "<display name> / <job> (<event>)"; the scoped one is the same but sets off the source repo with a colon,
+1 -1
View File
@@ -63,7 +63,7 @@ func handleInvalidWorkflows(ctx context.Context, input *notifyInput, ref git.Ref
continue
}
if err := createWorkflowCommitStatus(ctx, run.Repo, run.CommitSHA, entryName+" ("+run.TriggerEvent+")", run.WorkflowID,
commitstatus.CommitStatusFailure, run.Link(), "Invalid workflow file"); err != nil {
commitstatus.CommitStatusFailure, run.Link(), "Invalid workflow file", false); err != nil {
log.Error("create commit status for invalid workflow %q: %v", entryName, err)
}
NotifyWorkflowRunStatusUpdate(ctx, run)
+30 -5
View File
@@ -244,9 +244,9 @@ func checkRunConcurrency(ctx context.Context, run *actions_model.ActionRun) (*jo
return result, nil
}
// checkJobsOfCurrentRunAttempt resolves blocked jobs of the run's latest attempt.
// checkJobsOfCurrentRunAttempt resolves pending and blocked jobs of the run's latest attempt.
func checkJobsOfCurrentRunAttempt(ctx context.Context, run *actions_model.ActionRun) (*jobsCheckResult, error) {
// Approval is the only transition allowed to release an approval-pending run.
// Approval is the only transition allowed to release a run awaiting approval.
if run.NeedApproval {
return &jobsCheckResult{}, nil
}
@@ -366,7 +366,7 @@ func checkJobsOfCurrentRunAttempt(ctx context.Context, run *actions_model.Action
}
result.UpdatedJobs = append(result.UpdatedJobs, resolver.matrixUpdatedJobs...)
// Caller and matrix expansion both insert Blocked jobs, which only a follow-up pass resolves.
// Caller and matrix expansion both insert Pending or Blocked jobs, which only a follow-up pass resolves.
// Like the caller's children, matrix siblings are left out of result.Jobs and picked up there.
if expandedAnyCaller || resolver.matrixChanged {
result.RunIDsToReEmit = append(result.RunIDsToReEmit, run.ID)
@@ -398,7 +398,7 @@ func cancelFailedMatrixSiblings(ctx context.Context, jobs actions_model.ActionJo
type jobStatusResolver struct {
statuses map[int64]actions_model.Status
// sortedIDs are the keys of statuses, so blocked jobs are resolved in insertion order.
// sortedIDs are the keys of statuses, so jobs are resolved in insertion order.
// Resolve only ever rewrites statuses values, never its key set.
sortedIDs []int64
needs map[int64][]int64
@@ -498,6 +498,20 @@ func (r *jobStatusResolver) resolveCheckNeeds(id int64) (allDone, allSucceed boo
return allDone, allSucceed
}
func (r *jobStatusResolver) updateStatus(ctx context.Context, job *actions_model.ActionRunJob, status actions_model.Status) error {
cond := builder.Eq{"status": job.Status}
job.Status = status
affected, err := actions_model.UpdateRunJob(ctx, job, cond, "status")
if err != nil {
return err
}
if affected != 1 {
return fmt.Errorf("no affected for updating job %d", job.ID)
}
r.statuses[job.ID] = status
return nil
}
func (r *jobStatusResolver) resolve(ctx context.Context) (map[int64]actions_model.Status, error) {
ret := map[int64]actions_model.Status{}
@@ -509,7 +523,7 @@ func (r *jobStatusResolver) resolve(ctx context.Context) (map[int64]actions_mode
for _, id := range r.sortedIDs {
status := r.statuses[id]
actionRunJob := r.jobMap[id]
if status != actions_model.StatusBlocked {
if !status.In(actions_model.StatusPending, actions_model.StatusBlocked) {
continue
}
// An expanded caller has been resolved in an earlier pass, skip.
@@ -522,7 +536,18 @@ func (r *jobStatusResolver) resolve(ctx context.Context) (map[int64]actions_mode
continue
}
}
// past the run's own holds, a job is Pending exactly while its needs are unfinished
allDone, allSucceed := r.resolveCheckNeeds(id)
var err error
switch {
case status.IsPending() && allDone:
err = r.updateStatus(ctx, actionRunJob, actions_model.StatusBlocked)
case status.IsBlocked() && !allDone:
err = r.updateStatus(ctx, actionRunJob, actions_model.StatusPending)
}
if err != nil {
return nil, err
}
if !allDone {
continue
}
+43 -32
View File
@@ -58,11 +58,11 @@ func Test_jobStatusResolver_Resolve(t *testing.T) {
},
},
{
name: "multiple blocked",
name: "multiple pending",
jobs: actions_model.ActionJobList{
{ID: 1, JobID: "1", Status: actions_model.StatusSuccess, Needs: []string{}},
{ID: 2, JobID: "2", Status: actions_model.StatusBlocked, Needs: []string{"1"}},
{ID: 3, JobID: "3", Status: actions_model.StatusBlocked, Needs: []string{"1"}},
{ID: 2, JobID: "2", Status: actions_model.StatusPending, Needs: []string{"1"}},
{ID: 3, JobID: "3", Status: actions_model.StatusPending, Needs: []string{"1"}},
},
want: map[int64]actions_model.Status{
2: actions_model.StatusWaiting,
@@ -70,11 +70,11 @@ func Test_jobStatusResolver_Resolve(t *testing.T) {
},
},
{
name: "chain blocked",
name: "chain pending",
jobs: actions_model.ActionJobList{
{ID: 1, JobID: "1", Status: actions_model.StatusFailure, Needs: []string{}},
{ID: 2, JobID: "2", Status: actions_model.StatusBlocked, Needs: []string{"1"}},
{ID: 3, JobID: "3", Status: actions_model.StatusBlocked, Needs: []string{"2"}},
{ID: 2, JobID: "2", Status: actions_model.StatusPending, Needs: []string{"1"}},
{ID: 3, JobID: "3", Status: actions_model.StatusPending, Needs: []string{"2"}},
},
want: map[int64]actions_model.Status{
2: actions_model.StatusSkipped,
@@ -93,9 +93,9 @@ func Test_jobStatusResolver_Resolve(t *testing.T) {
{
name: "loop need",
jobs: actions_model.ActionJobList{
{ID: 1, JobID: "1", Status: actions_model.StatusBlocked, Needs: []string{"3"}},
{ID: 2, JobID: "2", Status: actions_model.StatusBlocked, Needs: []string{"1"}},
{ID: 3, JobID: "3", Status: actions_model.StatusBlocked, Needs: []string{"2"}},
{ID: 1, JobID: "1", Status: actions_model.StatusPending, Needs: []string{"3"}},
{ID: 2, JobID: "2", Status: actions_model.StatusPending, Needs: []string{"1"}},
{ID: 3, JobID: "3", Status: actions_model.StatusPending, Needs: []string{"2"}},
},
want: map[int64]actions_model.Status{},
},
@@ -103,7 +103,7 @@ func Test_jobStatusResolver_Resolve(t *testing.T) {
name: "`if` is not empty and all jobs in `needs` completed successfully",
jobs: actions_model.ActionJobList{
{ID: 1, JobID: "job1", Status: actions_model.StatusSuccess, Needs: []string{}},
{ID: 2, JobID: "job2", Status: actions_model.StatusBlocked, Needs: []string{"job1"}, WorkflowPayload: []byte(
{ID: 2, JobID: "job2", Status: actions_model.StatusPending, Needs: []string{"job1"}, WorkflowPayload: []byte(
`
name: test
on: push
@@ -122,7 +122,7 @@ jobs:
name: "`if` is not empty and not all jobs in `needs` completed successfully",
jobs: actions_model.ActionJobList{
{ID: 1, JobID: "job1", Status: actions_model.StatusFailure, Needs: []string{}},
{ID: 2, JobID: "job2", Status: actions_model.StatusBlocked, Needs: []string{"job1"}, WorkflowPayload: []byte(
{ID: 2, JobID: "job2", Status: actions_model.StatusPending, Needs: []string{"job1"}, WorkflowPayload: []byte(
`
name: test
on: push
@@ -141,7 +141,7 @@ jobs:
name: "`if` is empty and not all jobs in `needs` completed successfully",
jobs: actions_model.ActionJobList{
{ID: 1, JobID: "job1", Status: actions_model.StatusFailure, Needs: []string{}},
{ID: 2, JobID: "job2", Status: actions_model.StatusBlocked, Needs: []string{"job1"}, WorkflowPayload: []byte(
{ID: 2, JobID: "job2", Status: actions_model.StatusPending, Needs: []string{"job1"}, WorkflowPayload: []byte(
`
name: test
on: push
@@ -239,7 +239,7 @@ jobs:
name: "`if` is empty and a failed need has continue-on-error",
jobs: actions_model.ActionJobList{
{ID: 1, JobID: "job1", Status: actions_model.StatusFailure, ContinueOnError: true, Needs: []string{}},
{ID: 2, JobID: "job2", Status: actions_model.StatusBlocked, Needs: []string{"job1"}, WorkflowPayload: []byte(
{ID: 2, JobID: "job2", Status: actions_model.StatusPending, Needs: []string{"job1"}, WorkflowPayload: []byte(
`
name: test
on: push
@@ -263,7 +263,7 @@ jobs:
},
jobs: actions_model.ActionJobList{
{ID: 1, JobID: "job1", Status: actions_model.StatusSuccess, Needs: []string{}},
{ID: 2, JobID: "job2", Status: actions_model.StatusBlocked, Needs: []string{"job1"}, WorkflowPayload: []byte(
{ID: 2, JobID: "job2", Status: actions_model.StatusPending, Needs: []string{"job1"}, WorkflowPayload: []byte(
`
on:
workflow_dispatch:
@@ -288,10 +288,13 @@ jobs:
for i, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
// Each subtest gets a unique RunID / RunAttemptID so jobs from different subtests don't bleed into each other's FindTaskNeeds queries
runID := int64(9001 + i)
attemptID := int64(9001 + i)
repoID := tt.jobs[0].RepoID
dbRun := &actions_model.ActionRun{RepoID: repoID, Index: int64(9001 + i)}
require.NoError(t, db.Insert(ctx, dbRun))
attempt := &actions_model.ActionRunAttempt{RepoID: repoID, RunID: dbRun.ID, Attempt: 1}
require.NoError(t, db.Insert(ctx, attempt))
runID, attemptID := dbRun.ID, attempt.ID
run := util.IfZero(tt.run, stubRun)
require.NoError(t, db.Insert(ctx, &actions_model.ActionRunAttempt{ID: attemptID, RepoID: 1, RunID: runID}))
// Insert each test job (letting the DB assign IDs) and remember the testID -> dbID mapping so we can translate the expected map.
idMap := make(map[int64]int64, len(tt.jobs))
@@ -302,9 +305,7 @@ jobs:
j.RunAttemptID = attemptID
j.Run = run
// The resolver evaluates Blocked jobs via evaluateJobIf, which needs a valid YAML payload;
// supply a minimal one when the case didn't.
if j.Status == actions_model.StatusBlocked && len(j.WorkflowPayload) == 0 {
if j.Status.In(actions_model.StatusPending, actions_model.StatusBlocked) && len(j.WorkflowPayload) == 0 {
j.WorkflowPayload = minimalWorkflowPayload(j.JobID)
}
@@ -764,34 +765,44 @@ func Test_maxParallelReusableCallerLifecycle(t *testing.T) {
// dependents honest: a round resolved after an insert would judge them against a job set that is
// missing the siblings. See Resolve for why that is wrong.
func Test_jobStatusResolverStopsAfterMatrixInsert(t *testing.T) {
require.NoError(t, unittest.PrepareTestDatabase())
ctx := t.Context()
// build (2) stands for the expanded anchor: it reaches a terminal status this round, which is
// what would let report (3) resolve in the next one.
newChain := func() actions_model.ActionJobList {
return actions_model.ActionJobList{
{ID: 1, JobID: "generate", Status: actions_model.StatusFailure, WorkflowPayload: minimalWorkflowPayload("generate")},
{ID: 2, JobID: "build", Status: actions_model.StatusBlocked, Needs: []string{"generate"}, WorkflowPayload: minimalWorkflowPayload("build")},
{ID: 3, JobID: "report", Status: actions_model.StatusBlocked, Needs: []string{"build"}, WorkflowPayload: minimalWorkflowPayload("report")},
newChain := func(index int64) actions_model.ActionJobList {
run := &actions_model.ActionRun{Index: index}
require.NoError(t, db.Insert(ctx, run))
attempt := &actions_model.ActionRunAttempt{RunID: run.ID, Attempt: 1}
require.NoError(t, db.Insert(ctx, attempt))
jobs := actions_model.ActionJobList{
{JobID: "generate", Status: actions_model.StatusFailure, WorkflowPayload: minimalWorkflowPayload("generate")},
{JobID: "build", Status: actions_model.StatusPending, Needs: []string{"generate"}, WorkflowPayload: minimalWorkflowPayload("build")},
{JobID: "report", Status: actions_model.StatusPending, Needs: []string{"build"}, WorkflowPayload: minimalWorkflowPayload("report")},
}
for _, job := range jobs {
job.RunID, job.RunAttemptID = run.ID, attempt.ID
require.NoError(t, db.Insert(ctx, job))
}
return jobs
}
t.Run("without an insert the whole chain resolves in one pass", func(t *testing.T) {
got, err := newJobStatusResolver(newChain(), nil).Resolve(ctx)
chain := newChain(9201)
got, err := newJobStatusResolver(chain, nil).Resolve(ctx)
require.NoError(t, err)
assert.Equal(t, map[int64]actions_model.Status{
2: actions_model.StatusSkipped,
3: actions_model.StatusSkipped,
chain[1].ID: actions_model.StatusSkipped,
chain[2].ID: actions_model.StatusSkipped,
}, got)
})
t.Run("an insert stops the pass before the dependents are resolved", func(t *testing.T) {
r := newJobStatusResolver(newChain(), nil)
chain := newChain(9202)
r := newJobStatusResolver(chain, nil)
r.matrixInserted = true // as resolve() sets it once expansion has inserted siblings
got, err := r.Resolve(ctx)
require.NoError(t, err)
assert.Equal(t, map[int64]actions_model.Status{2: actions_model.StatusSkipped}, got,
assert.Equal(t, map[int64]actions_model.Status{chain[1].ID: actions_model.StatusSkipped}, got,
"report must wait for the re-emit, which sees the sibling combinations too")
})
}
+47
View File
@@ -4,6 +4,7 @@
package actions
import (
"slices"
"testing"
actions_model "gitea.dev/models/actions"
@@ -103,6 +104,52 @@ func TestPrepareRunAndInsert_MaxParallel(t *testing.T) {
}
}
func TestCheckJobs_ApprovedRunWaitsOnNeedsThenMaxParallel(t *testing.T) {
assert.NoError(t, unittest.PrepareTestDatabase())
defer test.MockVariableValue(&EmitJobsIfReadyByRun, func(int64) error { return nil })()
run := insertMaxParallelRun(t, `name: max-parallel
on: push
jobs:
setup:
runs-on: ubuntu-latest
steps:
- run: echo hi
build:
needs: setup
runs-on: ubuntu-latest
strategy:
max-parallel: 2
matrix:
version: [1, 2, 3]
steps:
- run: echo hi
`, true)
assert.Equal(t, map[actions_model.Status]int{actions_model.StatusBlocked: 4}, statusCounts(runJobs(t, run.ID, run.LatestAttemptID)))
repo := unittest.AssertExistsAndLoadBean(t, &repo_model.Repository{ID: run.RepoID})
doer := unittest.AssertExistsAndLoadBean(t, &user_model.User{ID: 1})
_, err := ApproveRuns(t.Context(), repo, doer, []int64{run.ID})
require.NoError(t, err)
checkJobs := func() actions_model.ActionJobList {
_, err := checkJobsOfCurrentRunAttempt(t.Context(), unittest.AssertExistsAndLoadBean(t, &actions_model.ActionRun{ID: run.ID}))
require.NoError(t, err)
return runJobs(t, run.ID, run.LatestAttemptID)
}
jobs := checkJobs()
assert.Equal(t, map[actions_model.Status]int{actions_model.StatusWaiting: 1, actions_model.StatusPending: 3}, statusCounts(jobs))
setup := jobs[slices.IndexFunc(jobs, func(job *actions_model.ActionRunJob) bool { return job.JobID == "setup" })]
setup.Status = actions_model.StatusSuccess
_, err = actions_model.UpdateRunJob(t.Context(), setup, nil, "status")
require.NoError(t, err)
assert.Equal(t, map[actions_model.Status]int{
actions_model.StatusSuccess: 1,
actions_model.StatusWaiting: 2,
actions_model.StatusBlocked: 1,
}, statusCounts(checkJobs()))
}
// A reusable workflow declares its own strategy, so the limit must reach the child jobs.
func TestInsertCallerChildren_MaxParallel(t *testing.T) {
assert.NoError(t, unittest.PrepareTestDatabase())
+1 -1
View File
@@ -28,7 +28,6 @@ func NotifyWorkflowJobsAndRunsStatusUpdate(ctx context.Context, jobs []*actions_
log.Error("Failed to load job attributes: %v", err)
continue
}
CreateCommitStatusForRunJobs(ctx, job.Run, job)
runRepoIDs[job.RunID] = job.RepoID
if _, ok := jobsByRunID[job.RunID]; !ok {
@@ -42,6 +41,7 @@ func NotifyWorkflowJobsAndRunsStatusUpdate(ctx context.Context, jobs []*actions_
}
for _, jobs := range jobsByRunID {
CreateCommitStatusForRunJobs(ctx, jobs[0].Run, jobs...)
NotifyWorkflowJobsStatusUpdate(ctx, jobs...)
}
}
+4 -3
View File
@@ -286,9 +286,10 @@ func execRerunPlan(ctx context.Context, plan *rerunPlan) (*actions_model.ActionR
var invalidIf error
if plan.rerunAttemptJobIDs.Contains(templateJob.AttemptJobID) {
// the emitter decides `if:` once all needs have results, and is the only place expanding a deferred matrix
shouldBlockJob := shouldBlock || len(newJob.Needs) > 0 || newJob.IsMatrixDeferred
newJob.Status = util.Iif(shouldBlockJob, actions_model.StatusBlocked, actions_model.StatusWaiting)
newJob.Status = util.Iif(shouldBlock, actions_model.StatusBlocked, actions_model.StatusWaiting)
if newJob.Status.IsWaiting() && len(newJob.Needs) > 0 {
newJob.Status = actions_model.StatusPending
}
newJob.TaskID = 0
newJob.SourceTaskID = 0
newJob.Started = 0
+3 -2
View File
@@ -25,6 +25,7 @@ import (
"gitea.dev/modules/log"
"gitea.dev/modules/setting"
api "gitea.dev/modules/structs"
"gitea.dev/modules/util"
"gitea.dev/services/convert"
"xorm.io/builder"
@@ -187,7 +188,7 @@ func canonicalCallUses(job *actions_model.ActionRunJob) string {
}
// expandReusableWorkflowCaller loads and parses the target reusable workflow and inserts the caller's direct child jobs.
// It expands only ONE level: a child that is itself a reusable caller is inserted Blocked and expanded later by a subsequent resolver pass.
// It expands only ONE level: a child that is itself a reusable caller is inserted Blocked or Pending and expanded later by a subsequent resolver pass.
// It does NOT schedule a follow-up resolver pass; the caller of this function is responsible for emitting.
//
// All call sites (PrepareRunAndInsert, execRerunPlan, checkJobsOfCurrentRunAttempt, ApproveRuns) invoke this inside their enclosing write transaction,
@@ -404,7 +405,7 @@ func insertCallerChildren(ctx context.Context, run *actions_model.ActionRun, att
RunsOn: parsedChild.RunsOn(),
ContinueOnError: parsedChild.GetContinueOnError(),
MaxParallel: parseMaxParallel(jobID, parsedChild.Strategy.MaxParallelString),
Status: actions_model.StatusBlocked,
Status: util.Iif(len(needs) > 0, actions_model.StatusPending, actions_model.StatusBlocked),
ParentJobID: caller.ID,
WorkflowSourceRepoID: sourceRepoID,
WorkflowSourceCommitSHA: sourceCommitSHA,
+6 -5
View File
@@ -191,7 +191,10 @@ func insertRunJob(ctx context.Context, run *actions_model.ActionRun, runAttempt
payload, _ := workflowJob.Marshal()
isReusableWorkflowCaller := job.Uses != ""
shouldBlockJob := runAttempt.Status == actions_model.StatusBlocked || len(needs) > 0 || run.NeedApproval
status := util.Iif(runAttempt.Status == actions_model.StatusBlocked || run.NeedApproval, actions_model.StatusBlocked, actions_model.StatusWaiting)
if status.IsWaiting() && len(needs) > 0 {
status = actions_model.StatusPending
}
attemptJobID, err := actions_model.GetNextAttemptJobID(ctx, run.ID)
if err != nil {
@@ -213,7 +216,7 @@ func insertRunJob(ctx context.Context, run *actions_model.ActionRun, runAttempt
AttemptJobID: attemptJobID,
Needs: needs,
RunsOn: job.RunsOn(),
Status: util.Iif(shouldBlockJob, actions_model.StatusBlocked, actions_model.StatusWaiting),
Status: status,
WorkflowSourceRepoID: run.WorkflowRepoID,
WorkflowSourceCommitSHA: run.WorkflowCommitSHA,
ContinueOnError: job.GetContinueOnError(),
@@ -256,9 +259,7 @@ func insertRunJob(ctx context.Context, run *actions_model.ActionRun, runAttempt
}
}
// If a job needs other jobs ("needs" is not empty), its status is set to StatusBlocked at the entry of the loop
// No need to check job concurrency for a blocked job (it will be checked by job emitter later)
// A slot-starved job skips the check too: it will not start, so it must not cancel its group peers.
// A slot-starved job skips the check: it will not start, so it must not cancel its group peers.
if runJob.Status == actions_model.StatusWaiting && slots.available(runJob) {
var jobsToCancel []*actions_model.ActionRunJob
runJob.Status, jobsToCancel, err = PrepareToStartJobWithConcurrency(ctx, runJob)