diff --git a/services/actions/job_emitter.go b/services/actions/job_emitter.go index 0b974b10de8..114e9d75ab3 100644 --- a/services/actions/job_emitter.go +++ b/services/actions/job_emitter.go @@ -362,15 +362,50 @@ func checkJobsOfCurrentRunAttempt(ctx context.Context, run *actions_model.Action } result.UpdatedJobs = append(result.UpdatedJobs, resolver.matrixUpdatedJobs...) - // Caller and matrix expansion insert Pending or Blocked jobs and a deferred gate leaves a job Blocked, only a follow-up pass resolves them. + hasFinishedCaller, err := reloadCallersFinishedByChildren(ctx, run.ID, resolver.jobMap, result.UpdatedJobs) + if err != nil { + return nil, err + } + // Only a follow-up pass resolves: + // - the children a caller expansion inserted, or the dependents of a caller that failed to expand + // - the siblings a matrix expansion inserted and the expanded job's dependents, or the dependents of a placeholder that failed to expand + // - a job the deferred gate left Blocked + // - the dependents of a caller finished by its children, which was still unfinished to the resolver // Like the caller's children, matrix siblings are left out of result.Jobs and picked up there. - if expandedAnyCaller || resolver.matrixChanged || resolver.gateDeferred { + if expandedAnyCaller || resolver.matrixChanged || resolver.gateDeferred || hasFinishedCaller { result.RunIDsToReEmit = append(result.RunIDsToReEmit, run.ID) } result.CancelledJobs = append(result.CancelledJobs, resolver.cancelledJobs...) return result, nil } +// reloadCallersFinishedByChildren reloads the callers the children's cascade finished in this pass, +// so this pass's commit statuses and run notification see the callers finished. +// They are kept out of result.UpdatedJobs, so no workflow_job webhook is sent for them, as for callers finished through a runner. +func reloadCallersFinishedByChildren(ctx context.Context, runID int64, jobMap map[int64]*actions_model.ActionRunJob, updatedJobs []*actions_model.ActionRunJob) (bool, error) { + hasFinishedCaller := false + checkedCallers := make(container.Set[int64]) + finishedJobs := slices.Clone(updatedJobs) + for i := 0; i < len(finishedJobs); i++ { + child := finishedJobs[i] + caller := jobMap[child.ParentJobID] + if caller == nil || !child.Status.IsDone() || caller.Status.IsDone() || !checkedCallers.Add(caller.ID) { + continue + } + freshCaller, err := actions_model.GetRunJobByRunAndID(ctx, runID, caller.ID) + if err != nil { + return false, fmt.Errorf("reloadCallersFinishedByChildren: reload caller %d: %w", caller.ID, err) + } + if !freshCaller.Status.IsDone() { + continue + } + caller.Status, caller.Started, caller.Stopped = freshCaller.Status, freshCaller.Started, freshCaller.Stopped + finishedJobs = append(finishedJobs, caller) + hasFinishedCaller = true + } + return hasFinishedCaller, nil +} + func cancelFailedMatrixSiblings(ctx context.Context, jobs actions_model.ActionJobList) ([]*actions_model.ActionRunJob, error) { var toCancel []*actions_model.ActionRunJob for _, failed := range jobs { diff --git a/services/actions/job_emitter_test.go b/services/actions/job_emitter_test.go index 6c96abab105..87702937cc5 100644 --- a/services/actions/job_emitter_test.go +++ b/services/actions/job_emitter_test.go @@ -5,6 +5,7 @@ package actions import ( "fmt" + "slices" "testing" actions_model "gitea.dev/models/actions" @@ -661,6 +662,65 @@ func Test_checkJobsOfCurrentRunAttempt_SkippedCallerIsUpdated(t *testing.T) { require.Len(t, result.UpdatedJobs, 1) assert.Equal(t, caller.ID, result.UpdatedJobs[0].ID) assert.Equal(t, actions_model.StatusSkipped, unittest.AssertExistsAndLoadBean(t, &actions_model.ActionRunJob{ID: caller.ID}).Status) + + // A child skipped in this pass finishes its caller and the outer caller, both are reported Success and the outer one's dependent is resolved by the re-emit. + run2 := &actions_model.ActionRun{ + RepoID: 4, OwnerID: 1, TriggerUserID: 1, + WorkflowID: "test.yml", Index: 9915, Ref: "refs/heads/main", Status: actions_model.StatusRunning, + } + assert.NoError(t, db.Insert(ctx, run2)) + attempt2 := &actions_model.ActionRunAttempt{RepoID: 4, RunID: run2.ID, Attempt: 1, Status: actions_model.StatusRunning} + assert.NoError(t, db.Insert(ctx, attempt2)) + _, err = db.Exec(ctx, "UPDATE `action_run` SET latest_attempt_id = ? WHERE id = ?", attempt2.ID, run2.ID) + assert.NoError(t, err) + run2.LatestAttemptID = attempt2.ID + caller2 := &actions_model.ActionRunJob{ + RunID: run2.ID, RunAttemptID: attempt2.ID, RepoID: 4, OwnerID: 1, + JobID: "deploy", Name: "deploy", Status: actions_model.StatusRunning, IsReusableCaller: true, IsExpanded: true, + WorkflowPayload: []byte("jobs: {deploy: {uses: ./.gitea/workflows/mid.yml}}"), + ReusableWorkflowContent: []byte("on: {workflow_call: {}}\njobs: {inner: {uses: ./.gitea/workflows/leaf.yml}}"), + } + assert.NoError(t, db.Insert(ctx, caller2)) + inner := &actions_model.ActionRunJob{ + RunID: run2.ID, RunAttemptID: attempt2.ID, RepoID: 4, OwnerID: 1, ParentJobID: caller2.ID, + JobID: "inner", Name: "inner", Status: actions_model.StatusRunning, IsReusableCaller: true, IsExpanded: true, + WorkflowPayload: []byte("jobs: {inner: {uses: ./.gitea/workflows/leaf.yml}}"), + } + assert.NoError(t, db.Insert(ctx, inner)) + assert.NoError(t, db.Insert(ctx, &actions_model.ActionRunJob{ + RunID: run2.ID, RunAttemptID: attempt2.ID, RepoID: 4, OwnerID: 1, ParentJobID: inner.ID, + JobID: "work", Status: actions_model.StatusSuccess, WorkflowPayload: minimalWorkflowPayload("work"), + })) + assert.NoError(t, db.Insert(ctx, &actions_model.ActionRunJob{ + RunID: run2.ID, RunAttemptID: attempt2.ID, RepoID: 4, OwnerID: 1, ParentJobID: inner.ID, + JobID: "alert", Status: actions_model.StatusBlocked, Needs: []string{"work"}, + WorkflowPayload: []byte("jobs: {alert: {if: false, needs: [work], runs-on: ubuntu-latest, steps: [{run: echo}]}}"), + })) + after := &actions_model.ActionRunJob{ + RunID: run2.ID, RunAttemptID: attempt2.ID, RepoID: 4, OwnerID: 1, + JobID: "after", Status: actions_model.StatusPending, Needs: []string{"deploy"}, + WorkflowPayload: []byte("jobs: {after: {if: \"github.ref == 'refs/heads/main'\", needs: [deploy], runs-on: ubuntu-latest, steps: [{run: echo}]}}"), + } + assert.NoError(t, db.Insert(ctx, after)) + + result, err = checkJobsOfCurrentRunAttempt(ctx, run2) + assert.NoError(t, err) + assert.Contains(t, result.RunIDsToReEmit, run2.ID) + assert.Equal(t, actions_model.StatusPending, unittest.AssertExistsAndLoadBean(t, &actions_model.ActionRunJob{ID: after.ID}).Status) + updatedJobIDs := make([]string, 0, len(result.UpdatedJobs)) + for _, job := range result.UpdatedJobs { + updatedJobIDs = append(updatedJobIDs, job.JobID) + } + assert.Equal(t, []string{"alert"}, updatedJobIDs) + for _, callerID := range []int64{inner.ID, caller2.ID} { + idx := slices.IndexFunc(result.Jobs, func(job *actions_model.ActionRunJob) bool { return job.ID == callerID }) + require.NotEqual(t, -1, idx) + assert.Equal(t, actions_model.StatusSuccess, result.Jobs[idx].Status) + } + + _, err = checkJobsOfCurrentRunAttempt(ctx, run2) + assert.NoError(t, err) + assert.Equal(t, actions_model.StatusWaiting, unittest.AssertExistsAndLoadBean(t, &actions_model.ActionRunJob{ID: after.ID}).Status) } // Test_checkRunConcurrency_HeldGroupDoesNotWake verifies that only an unoccupied concurrency group can wake up a blocked run/job.