mirror of
https://github.com/go-gitea/gitea.git
synced 2026-08-05 10:10:14 +00:00
11d0ed699b
Follow-up to https://github.com/go-gitea/gitea/pull/36564 (dynamic matrix) and https://github.com/go-gitea/gitea/pull/36357 (max-parallel), fixing issues found reviewing the two features together. - **A placeholder could stall its run forever.** Its payload keeps the raw matrix but loses its `needs`, so `ParseJob` re-expanded it instead of reading it back — fatal for `include: ${{ fromJson(needs.*.outputs.*) }}`. - **An `if:` reading `matrix.*` skipped the whole job**, with or without the `${{ }}`. It now reduces to the needs gate, except under `always()`/`failure()`/`cancelled()`, and each combination is decided on its own values once the matrix expands. - **Dependents could be skipped before the combinations ran**, since inserted siblings are absent from the resolver's job set. The pass now stops after an insert and defers to the re-emit it schedules. - **Expansion failures stranded the placeholder.** A retryable one is returned so the queue retries it; a malformed payload fails the job instead of requeueing forever. - **Rerun could rewind a pass-through row** into a raw placeholder keeping its old terminal status, which nothing expands. Now gated on the anchor itself. Plus: `max-parallel` distinguishes an unevaluated `${{ }}` (debug) from a non-numeric literal (warn — it silently drops the cap). Co-authored-by: Zettat123 <zettat123@gmail.com> Co-authored-by: silverwind <me@silverwind.io>
611 lines
21 KiB
Go
611 lines
21 KiB
Go
// Copyright 2022 The Gitea Authors. All rights reserved.
|
|
// SPDX-License-Identifier: MIT
|
|
|
|
package actions
|
|
|
|
import (
|
|
"context"
|
|
"errors"
|
|
"fmt"
|
|
"slices"
|
|
|
|
actions_model "gitea.dev/models/actions"
|
|
"gitea.dev/models/db"
|
|
"gitea.dev/modules/container"
|
|
"gitea.dev/modules/graceful"
|
|
"gitea.dev/modules/log"
|
|
"gitea.dev/modules/queue"
|
|
"gitea.dev/modules/setting"
|
|
"gitea.dev/modules/timeutil"
|
|
|
|
"xorm.io/builder"
|
|
)
|
|
|
|
var jobEmitterQueue *queue.WorkerPoolQueue[*jobUpdate]
|
|
|
|
type jobUpdate struct {
|
|
RunID int64
|
|
}
|
|
|
|
var EmitJobsIfReadyByRun = func(runID int64) error {
|
|
err := jobEmitterQueue.Push(&jobUpdate{
|
|
RunID: runID,
|
|
})
|
|
if errors.Is(err, queue.ErrAlreadyInQueue) {
|
|
return nil
|
|
}
|
|
return err
|
|
}
|
|
|
|
func EmitJobsIfReadyByJobs(jobs []*actions_model.ActionRunJob) {
|
|
checkedRuns := make(container.Set[int64])
|
|
for _, job := range jobs {
|
|
if !job.Status.IsDone() || checkedRuns.Contains(job.RunID) {
|
|
continue
|
|
}
|
|
if err := EmitJobsIfReadyByRun(job.RunID); err != nil {
|
|
log.Error("Check jobs of run %d: %v", job.RunID, err)
|
|
}
|
|
checkedRuns.Add(job.RunID)
|
|
}
|
|
}
|
|
|
|
func jobEmitterQueueHandler(items ...*jobUpdate) []*jobUpdate {
|
|
ctx := graceful.GetManager().ShutdownContext()
|
|
var ret []*jobUpdate
|
|
for _, update := range items {
|
|
if err := checkJobsByRunID(ctx, update.RunID); err != nil {
|
|
log.Error("check run %d: %v", update.RunID, err)
|
|
ret = append(ret, update)
|
|
}
|
|
}
|
|
return ret
|
|
}
|
|
|
|
func checkJobsByRunID(ctx context.Context, runID int64) error {
|
|
run, exist, err := db.GetByID[actions_model.ActionRun](ctx, runID)
|
|
if !exist {
|
|
return fmt.Errorf("run %d does not exist", runID)
|
|
}
|
|
if err != nil {
|
|
return fmt.Errorf("get action run: %w", err)
|
|
}
|
|
var result jobsCheckResult
|
|
if err := db.WithTx(ctx, func(ctx context.Context) error {
|
|
// check jobs of the current run
|
|
r, err := checkJobsOfCurrentRunAttempt(ctx, run)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
result.merge(r)
|
|
|
|
r, err = checkRunConcurrency(ctx, run)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
result.merge(r)
|
|
return nil
|
|
}); err != nil {
|
|
return err
|
|
}
|
|
// Re-emit AFTER the transaction commits; doing this inside WithTx would deadlock under
|
|
// immediate-mode queues (the inline handler reopens checkJobsByRunID and asks for a
|
|
// nested writer transaction while the outer one is still open).
|
|
emitted := make(container.Set[int64])
|
|
for _, rid := range result.RunIDsToReEmit {
|
|
if !emitted.Add(rid) {
|
|
continue
|
|
}
|
|
if err := EmitJobsIfReadyByRun(rid); err != nil {
|
|
log.Error("re-emit run %d: %v", rid, err)
|
|
}
|
|
}
|
|
NotifyWorkflowJobsAndRunsStatusUpdate(ctx, result.CancelledJobs)
|
|
EmitJobsIfReadyByJobs(result.CancelledJobs)
|
|
if err := createCommitStatusesForJobsByRun(ctx, result.Jobs); err != nil {
|
|
return err
|
|
}
|
|
NotifyWorkflowJobsStatusUpdate(ctx, result.UpdatedJobs...)
|
|
runJobs := make(map[int64][]*actions_model.ActionRunJob)
|
|
for _, job := range result.Jobs {
|
|
runJobs[job.RunID] = append(runJobs[job.RunID], job)
|
|
}
|
|
runUpdatedJobs := make(map[int64][]*actions_model.ActionRunJob)
|
|
for _, uj := range result.UpdatedJobs {
|
|
runUpdatedJobs[uj.RunID] = append(runUpdatedJobs[uj.RunID], uj)
|
|
}
|
|
for runID, js := range runJobs {
|
|
if len(runUpdatedJobs[runID]) == 0 {
|
|
continue
|
|
}
|
|
runUpdated := true
|
|
for _, job := range js {
|
|
if !job.Status.IsDone() {
|
|
runUpdated = false
|
|
break
|
|
}
|
|
}
|
|
if runUpdated {
|
|
NotifyWorkflowRunStatusUpdateWithReload(ctx, js[0].RepoID, js[0].RunID)
|
|
}
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func createCommitStatusesForJobsByRun(ctx context.Context, jobs []*actions_model.ActionRunJob) error {
|
|
runJobs := make(map[int64][]*actions_model.ActionRunJob)
|
|
for _, job := range jobs {
|
|
runJobs[job.RunID] = append(runJobs[job.RunID], job)
|
|
}
|
|
|
|
for jobRunID, jobList := range runJobs {
|
|
run, err := actions_model.GetRunByRepoAndID(ctx, jobList[0].RepoID, jobRunID)
|
|
if err != nil {
|
|
return fmt.Errorf("get action run %d: %w", jobRunID, err)
|
|
}
|
|
CreateCommitStatusForRunJobs(ctx, run, jobList...)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// findConcurrencyWaiterToWake returns a run (other than excludeRunID) blocked on the group that can be woken now
|
|
func findConcurrencyWaiterToWake(ctx context.Context, repoID, excludeRunID int64, concurrencyGroup string) (int64, error) {
|
|
if concurrencyGroup == "" {
|
|
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})
|
|
if err != nil {
|
|
return 0, fmt.Errorf("find blocked concurrent runs: %w", err)
|
|
}
|
|
for _, a := range cAttempts {
|
|
if a.RunID != excludeRunID {
|
|
return a.RunID, nil
|
|
}
|
|
}
|
|
for _, j := range cJobs {
|
|
if j.RunID != excludeRunID {
|
|
return j.RunID, nil
|
|
}
|
|
}
|
|
return 0, nil
|
|
}
|
|
|
|
// checkRunConcurrency wakes a run blocked by concurrency that may become runnable now that
|
|
// the current run's activity may have freed a workflow-level or job-level concurrency group.
|
|
func checkRunConcurrency(ctx context.Context, run *actions_model.ActionRun) (*jobsCheckResult, error) {
|
|
result := &jobsCheckResult{}
|
|
checkedConcurrencyGroup := make(container.Set[string])
|
|
|
|
collect := func(concurrencyGroup string) error {
|
|
checkedConcurrencyGroup.Add(concurrencyGroup)
|
|
|
|
// Exclude run.ID: this run's own jobs are resolved by checkJobsOfCurrentRunAttempt, no need to re-emit.
|
|
concurrentRunID, err := findConcurrencyWaiterToWake(ctx, run.RepoID, run.ID, concurrencyGroup)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if concurrentRunID == 0 {
|
|
return nil
|
|
}
|
|
|
|
concurrentRun, err := actions_model.GetRunByRepoAndID(ctx, run.RepoID, concurrentRunID)
|
|
if err != nil {
|
|
return fmt.Errorf("get concurrent run %d: %w", concurrentRunID, err)
|
|
}
|
|
// A run awaiting approval is not advanced by concurrency; ApproveRuns emits it once approved.
|
|
if concurrentRun.NeedApproval {
|
|
return nil
|
|
}
|
|
|
|
result.RunIDsToReEmit = append(result.RunIDsToReEmit, concurrentRunID)
|
|
return nil
|
|
}
|
|
|
|
// check run (workflow-level) concurrency
|
|
runConcurrencyGroup, _, err := run.GetEffectiveConcurrency(ctx)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("GetEffectiveConcurrency: %w", err)
|
|
}
|
|
if runConcurrencyGroup != "" {
|
|
if err := collect(runConcurrencyGroup); err != nil {
|
|
return nil, err
|
|
}
|
|
}
|
|
|
|
// check job concurrency
|
|
runJobs, err := actions_model.GetLatestAttemptJobsByRepoAndRunID(ctx, run.RepoID, run.ID)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("find run %d jobs: %w", run.ID, err)
|
|
}
|
|
for _, job := range runJobs {
|
|
if !job.Status.IsDone() {
|
|
continue
|
|
}
|
|
if job.ConcurrencyGroup == "" || checkedConcurrencyGroup.Contains(job.ConcurrencyGroup) {
|
|
continue
|
|
}
|
|
if err := collect(job.ConcurrencyGroup); err != nil {
|
|
return nil, err
|
|
}
|
|
}
|
|
return result, nil
|
|
}
|
|
|
|
// checkJobsOfCurrentRunAttempt resolves blocked jobs of the run's latest attempt.
|
|
func checkJobsOfCurrentRunAttempt(ctx context.Context, run *actions_model.ActionRun) (*jobsCheckResult, error) {
|
|
jobs, err := actions_model.GetRunJobsByRunAndAttemptID(ctx, run.ID, run.LatestAttemptID)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
result := &jobsCheckResult{Jobs: jobs}
|
|
|
|
var attempt *actions_model.ActionRunAttempt
|
|
if run.LatestAttemptID > 0 {
|
|
attempt, err = actions_model.GetRunAttemptByRepoAndID(ctx, run.RepoID, run.LatestAttemptID)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
}
|
|
|
|
// The resolver below only considers needs and job-level concurrency, so a run blocked
|
|
// solely by run-level concurrency would have its jobs unblocked here. checkRunConcurrency
|
|
// re-evaluates when the holding run finishes.
|
|
if run.Status.IsBlocked() && attempt != nil {
|
|
shouldBlock, err := shouldBlockRunByConcurrency(ctx, attempt)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("shouldBlockRunByConcurrency: %w", err)
|
|
}
|
|
if shouldBlock {
|
|
return result, nil
|
|
}
|
|
}
|
|
|
|
vars, err := actions_model.GetVariablesOfRun(ctx, run)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
resolver := newJobStatusResolver(jobs, vars)
|
|
|
|
expandedAnyCaller := false
|
|
if err = db.WithTx(ctx, func(ctx context.Context) error {
|
|
for _, job := range jobs {
|
|
job.Run = run
|
|
}
|
|
|
|
updates, err := resolver.Resolve(ctx)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
for _, job := range jobs {
|
|
status, ok := updates[job.ID]
|
|
if !ok {
|
|
continue
|
|
}
|
|
|
|
if job.IsReusableCaller {
|
|
switch status {
|
|
case actions_model.StatusWaiting:
|
|
if err := expandReusableWorkflowCaller(ctx, run, attempt, job, vars); err != nil {
|
|
// Terminal expansion failure (an invalid/unresolvable reusable workflow): fail this caller.
|
|
log.Warn("caller %d cannot be expanded: %v", job.ID, err)
|
|
job.Status = actions_model.StatusFailure
|
|
job.Stopped = timeutil.TimeStampNow()
|
|
if n, uerr := actions_model.UpdateRunJob(ctx, job, builder.Eq{"status": actions_model.StatusBlocked, "is_expanded": false}, "status", "stopped"); uerr != nil {
|
|
return fmt.Errorf("mark unexpandable caller %d failed: %w", job.ID, uerr)
|
|
} else if n == 1 {
|
|
log.Warn("unexpandable caller %d has been marked as failed", job.ID)
|
|
result.UpdatedJobs = append(result.UpdatedJobs, job)
|
|
// Re-emit so the failed caller's dependents get resolved on the next pass.
|
|
expandedAnyCaller = true
|
|
} else {
|
|
// A concurrent writer advanced the caller; restore the in-memory state.
|
|
log.Warn("unexpandable caller %d has been advanced by a concurrent writer, not marking it failed", job.ID)
|
|
job.Status = actions_model.StatusBlocked
|
|
job.Stopped = 0
|
|
}
|
|
} else {
|
|
expandedAnyCaller = true
|
|
}
|
|
case actions_model.StatusSkipped:
|
|
job.Status = actions_model.StatusSkipped
|
|
if _, err := actions_model.UpdateRunJob(ctx, job, nil, "status"); err != nil {
|
|
return err
|
|
}
|
|
}
|
|
continue
|
|
}
|
|
|
|
// Non-caller: standard status update.
|
|
job.Status = status
|
|
if n, err := actions_model.UpdateRunJob(ctx, job, builder.Eq{"status": actions_model.StatusBlocked}, "status"); err != nil {
|
|
return err
|
|
} else if n != 1 {
|
|
return fmt.Errorf("no affected for updating blocked job %v", job.ID)
|
|
}
|
|
result.UpdatedJobs = append(result.UpdatedJobs, job)
|
|
}
|
|
return nil
|
|
}); err != nil {
|
|
return nil, err
|
|
}
|
|
|
|
result.UpdatedJobs = append(result.UpdatedJobs, resolver.matrixUpdatedJobs...)
|
|
// Caller and matrix expansion both insert 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)
|
|
}
|
|
result.CancelledJobs = resolver.cancelledJobs
|
|
return result, nil
|
|
}
|
|
|
|
type jobStatusResolver struct {
|
|
statuses map[int64]actions_model.Status
|
|
// sortedIDs are the keys of statuses, so blocked jobs are resolved in insertion order.
|
|
// Resolve only ever rewrites statuses values, never its key set.
|
|
sortedIDs []int64
|
|
needs map[int64][]int64
|
|
jobMap map[int64]*actions_model.ActionRunJob
|
|
vars map[string]string
|
|
cancelledJobs []*actions_model.ActionRunJob
|
|
// matrixChanged is set when matrix expansion inserted siblings or failed a placeholder, both of
|
|
// which need a follow-up pass to resolve the dependents.
|
|
matrixChanged bool
|
|
// matrixInserted is set when matrix expansion inserted sibling rows, which Resolve stops on.
|
|
matrixInserted bool
|
|
// 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
|
|
}
|
|
|
|
func newJobStatusResolver(jobs actions_model.ActionJobList, vars map[string]string) *jobStatusResolver {
|
|
// Scope-aware: needs are resolved within the same ParentJobID scope so the same
|
|
// JobID in different reusable workflow calls does not cross-link.
|
|
scopedIDToJobs := make(map[int64]map[string][]*actions_model.ActionRunJob)
|
|
jobMap := make(map[int64]*actions_model.ActionRunJob)
|
|
for _, job := range jobs {
|
|
scope := scopedIDToJobs[job.ParentJobID]
|
|
if scope == nil {
|
|
scope = make(map[string][]*actions_model.ActionRunJob)
|
|
scopedIDToJobs[job.ParentJobID] = scope
|
|
}
|
|
scope[job.JobID] = append(scope[job.JobID], job)
|
|
jobMap[job.ID] = job
|
|
}
|
|
|
|
statuses := make(map[int64]actions_model.Status, len(jobs))
|
|
needs := make(map[int64][]int64, len(jobs))
|
|
sortedIDs := make([]int64, 0, len(jobs))
|
|
for _, job := range jobs {
|
|
statuses[job.ID] = job.Status
|
|
sortedIDs = append(sortedIDs, job.ID)
|
|
scope := scopedIDToJobs[job.ParentJobID]
|
|
for _, need := range job.Needs {
|
|
for _, v := range scope[need] {
|
|
needs[job.ID] = append(needs[job.ID], v.ID)
|
|
}
|
|
}
|
|
}
|
|
slices.Sort(sortedIDs)
|
|
return &jobStatusResolver{
|
|
statuses: statuses,
|
|
sortedIDs: sortedIDs,
|
|
needs: needs,
|
|
jobMap: jobMap,
|
|
vars: vars,
|
|
}
|
|
}
|
|
|
|
func (r *jobStatusResolver) Resolve(ctx context.Context) (map[int64]actions_model.Status, error) {
|
|
ret := map[int64]actions_model.Status{}
|
|
for i := 0; i < len(r.statuses); i++ {
|
|
updated, err := r.resolve(ctx)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
if len(updated) == 0 {
|
|
return ret, nil
|
|
}
|
|
for k, v := range updated {
|
|
ret[k] = v
|
|
r.statuses[k] = v
|
|
}
|
|
if r.matrixInserted {
|
|
// Matrix expansion inserted sibling rows this round. They are not in statuses/needs, so
|
|
// another round would resolve a dependent of the expanded job against the placeholder's
|
|
// own combination alone: if that combination was just skipped by its `if:`, the dependent
|
|
// sees all its needs done and gets skipped before any sibling has even started. Stop here
|
|
// and let the re-emit, which reloads the full job set, resolve them.
|
|
return ret, nil
|
|
}
|
|
}
|
|
return ret, nil
|
|
}
|
|
|
|
func (r *jobStatusResolver) resolveCheckNeeds(id int64) (allDone, allSucceed bool) {
|
|
allDone, allSucceed = true, true
|
|
for _, need := range r.needs[id] {
|
|
needStatus := r.statuses[need]
|
|
if !needStatus.IsDone() {
|
|
allDone = false
|
|
}
|
|
// A failed need with continue-on-error:true is treated as success, matching AggregateJobStatus,
|
|
// so a downstream job with an implicit `success()` is not skipped.
|
|
if needJob := r.jobMap[need]; needJob != nil && needJob.ContinueOnError && needStatus == actions_model.StatusFailure {
|
|
continue
|
|
}
|
|
if needStatus.In(actions_model.StatusFailure, actions_model.StatusCancelled, actions_model.StatusSkipped) {
|
|
allSucceed = false
|
|
}
|
|
}
|
|
return allDone, allSucceed
|
|
}
|
|
|
|
func (r *jobStatusResolver) resolve(ctx context.Context) (map[int64]actions_model.Status, error) {
|
|
ret := map[int64]actions_model.Status{}
|
|
|
|
slots := maxParallelSlots{}
|
|
for id, status := range r.statuses {
|
|
slots.hold(r.jobMap[id], status)
|
|
}
|
|
|
|
for _, id := range r.sortedIDs {
|
|
status := r.statuses[id]
|
|
actionRunJob := r.jobMap[id]
|
|
if status != actions_model.StatusBlocked {
|
|
continue
|
|
}
|
|
// An expanded caller has been resolved in an earlier pass, skip.
|
|
if actionRunJob.IsReusableCaller && actionRunJob.IsExpanded {
|
|
continue
|
|
}
|
|
// A child of a caller cannot start until the caller has become "ready" (children inserted, CallPayload populated).
|
|
if actionRunJob.ParentJobID > 0 {
|
|
if parent, ok := r.jobMap[actionRunJob.ParentJobID]; ok && !parent.IsExpanded {
|
|
continue
|
|
}
|
|
}
|
|
allDone, allSucceed := r.resolveCheckNeeds(id)
|
|
if !allDone {
|
|
continue
|
|
}
|
|
|
|
// Decide whether the job runs at all before expanding a deferred matrix: a job whose needs
|
|
// failed or were skipped has to be skipped too, not failed for a matrix those needs never
|
|
// produced the outputs for. An `if:` that reads `matrix.*` cannot be decided this early, so
|
|
// evaluateJobIf reduces it to that needs gate and the pass below decides it per combination.
|
|
shouldStartJob, err := evaluateJobIf(ctx, actionRunJob.Run, nil, actionRunJob, r.vars, allSucceed)
|
|
if err != nil {
|
|
// TODO: surface deterministic expression errors to users by failing the job with a message.
|
|
log.Error("evaluateJobIf failed, job will stay blocked: job: %d, err: %v", id, err)
|
|
continue
|
|
}
|
|
if !shouldStartJob {
|
|
ret[id] = actions_model.StatusSkipped
|
|
continue
|
|
}
|
|
|
|
// Expand a needs-dependent matrix now that its needs are done and the job is going to run.
|
|
wasDeferred := actionRunJob.IsMatrixDeferred
|
|
siblings, err := expandDeferredMatrix(ctx, actionRunJob, r.vars)
|
|
if err != nil {
|
|
// Aborting the pass is required: once the placeholder is claimed as the first combination,
|
|
// committing here would drop the remaining ones for good. Before the claim it is what gets
|
|
// the pass retried by the job-emitter queue, since a run whose needs are all done has
|
|
// nothing left to trigger another pass on its own.
|
|
return nil, fmt.Errorf("expand matrix of job %d: %w", id, err)
|
|
}
|
|
if actionRunJob.Status != actions_model.StatusBlocked {
|
|
// expandDeferredMatrix already persisted the failure, so it bypasses `ret`.
|
|
r.statuses[id] = actionRunJob.Status
|
|
r.matrixUpdatedJobs = append(r.matrixUpdatedJobs, actionRunJob)
|
|
r.matrixChanged = true
|
|
continue
|
|
}
|
|
if actionRunJob.IsMatrixDeferred {
|
|
continue // could not be expanded yet, it stays blocked and is retried on the next pass
|
|
}
|
|
if len(siblings) > 0 {
|
|
r.matrixChanged, r.matrixInserted = true, true
|
|
}
|
|
if wasDeferred {
|
|
// This row is now the first combination, and the `if:` can be evaluated.
|
|
// Gate it on its own combination here, as the siblings will be on the next pass.
|
|
shouldStartJob, err := evaluateJobIf(ctx, actionRunJob.Run, nil, actionRunJob, r.vars, allSucceed)
|
|
if err != nil {
|
|
log.Error("evaluateJobIf failed after matrix expansion, job will stay blocked: job: %d, err: %v", id, err)
|
|
continue
|
|
}
|
|
if !shouldStartJob {
|
|
ret[id] = actions_model.StatusSkipped
|
|
continue
|
|
}
|
|
}
|
|
|
|
// A slot-starved job cannot start, skip the following checks.
|
|
if !slots.available(actionRunJob) {
|
|
continue
|
|
}
|
|
|
|
// update concurrency and check whether the job can run now
|
|
err = updateConcurrencyEvaluationForJobWithNeeds(ctx, actionRunJob, r.vars)
|
|
if err != nil {
|
|
// The err can be caused by different cases: database error, or syntax error, or the needed jobs haven't completed
|
|
// At the moment there is no way to distinguish them.
|
|
// TODO: if workflow or concurrency expression has syntax error, there should be a user error message, need to show it to end users
|
|
log.Debug("updateConcurrencyEvaluationForJobWithNeeds failed, this job will stay blocked: job: %d, err: %v", id, err)
|
|
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...)
|
|
}
|
|
|
|
if newStatus == actions_model.StatusWaiting && !slots.take(actionRunJob) {
|
|
continue // no free slot, leave blocked
|
|
}
|
|
|
|
if newStatus != actions_model.StatusBlocked {
|
|
ret[id] = newStatus
|
|
}
|
|
}
|
|
return ret, nil
|
|
}
|
|
|
|
func updateConcurrencyEvaluationForJobWithNeeds(ctx context.Context, actionRunJob *actions_model.ActionRunJob, vars map[string]string) error {
|
|
if setting.IsInTesting && actionRunJob.RepoID == 0 {
|
|
return nil // for testing purpose only, no repo, no evaluation
|
|
}
|
|
|
|
// Legacy jobs (created before migration v331) have RunAttemptID=0 and no attempt record.
|
|
var attempt *actions_model.ActionRunAttempt
|
|
if actionRunJob.RunAttemptID > 0 {
|
|
var err error
|
|
attempt, err = actions_model.GetRunAttemptByRepoAndID(ctx, actionRunJob.RepoID, actionRunJob.RunAttemptID)
|
|
if err != nil {
|
|
return fmt.Errorf("GetRunAttemptByRepoAndID: %w", err)
|
|
}
|
|
}
|
|
if err := EvaluateJobConcurrencyFillModel(ctx, actionRunJob.Run, attempt, actionRunJob, vars, nil); err != nil {
|
|
return fmt.Errorf("evaluate job concurrency: %w", err)
|
|
}
|
|
|
|
if _, err := actions_model.UpdateRunJob(ctx, actionRunJob, nil, "concurrency_group", "concurrency_cancel", "is_concurrency_evaluated"); err != nil {
|
|
return fmt.Errorf("update run job: %w", err)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// jobsCheckResult bundles the output of the per-run job-check helpers.
|
|
type jobsCheckResult struct {
|
|
// Jobs are all jobs of the run's latest attempt that were inspected.
|
|
Jobs []*actions_model.ActionRunJob
|
|
// UpdatedJobs are jobs whose status was transitioned out of Blocked in this pass.
|
|
UpdatedJobs []*actions_model.ActionRunJob
|
|
// CancelledJobs are jobs cancelled by job-level concurrency while preparing to start.
|
|
CancelledJobs []*actions_model.ActionRunJob
|
|
// RunIDsToReEmit are runs that need another resolver pass in their own transaction.
|
|
RunIDsToReEmit []int64
|
|
}
|
|
|
|
// merge appends another result's contents into r in place.
|
|
func (r *jobsCheckResult) merge(other *jobsCheckResult) {
|
|
r.Jobs = append(r.Jobs, other.Jobs...)
|
|
r.UpdatedJobs = append(r.UpdatedJobs, other.UpdatedJobs...)
|
|
r.CancelledJobs = append(r.CancelledJobs, other.CancelledJobs...)
|
|
r.RunIDsToReEmit = append(r.RunIDsToReEmit, other.RunIDsToReEmit...)
|
|
}
|