mirror of
https://github.com/go-gitea/gitea.git
synced 2026-09-27 18:24:21 +00:00
3875db1974
Adds a read-only Actions job queue: running jobs first, then waiting jobs in the order a runner picks them up. It is shown instance-wide in the admin Actions section with owner, repository and status filters, and per repository in the Actions tab. Both lists refresh in place. Pending work is currently only visible per repository and newest-first, so nothing shows what is queued, in which order, or what occupies the runners. Reordering the queue will be proposed separately. A migration adds indexes for the runner pickup query and repository-scoped status lookups. * Fix #34198 <img width="1345" height="451" alt="image" src="https://github.com/user-attachments/assets/7d52ff76-81b4-44e8-b583-d7d89c9dffcd" /> <img width="1809" height="1134" alt="image" src="https://github.com/user-attachments/assets/4d56c0cb-bae7-4ce2-8f3c-75163b2bc7f4" /> --------- Co-authored-by: Zettat123 <zettat123@gmail.com> Co-authored-by: silverwind <me@silverwind.io> Co-authored-by: wxiaoguang <wxiaoguang@gmail.com>
671 lines
20 KiB
Go
671 lines
20 KiB
Go
// Copyright 2022 The Gitea Authors. All rights reserved.
|
|
// SPDX-License-Identifier: MIT
|
|
|
|
package actions
|
|
|
|
import (
|
|
"context"
|
|
"errors"
|
|
"fmt"
|
|
"strings"
|
|
"time"
|
|
|
|
runnerv1 "gitea.dev/actionslib/runner/v1"
|
|
auth_model "gitea.dev/models/auth"
|
|
"gitea.dev/models/db"
|
|
"gitea.dev/models/unit"
|
|
"gitea.dev/modules/actions/jobparser"
|
|
"gitea.dev/modules/globallock"
|
|
"gitea.dev/modules/log"
|
|
"gitea.dev/modules/setting"
|
|
"gitea.dev/modules/timeutil"
|
|
"gitea.dev/modules/util"
|
|
|
|
"xorm.io/builder"
|
|
)
|
|
|
|
// ActionTask represents a distribution of job
|
|
type ActionTask struct {
|
|
ID int64
|
|
JobID int64
|
|
Job *ActionRunJob `xorm:"-"`
|
|
Steps []*ActionTaskStep `xorm:"-"`
|
|
Attempt int64
|
|
RunnerID int64 `xorm:"index"`
|
|
Status Status `xorm:"index"`
|
|
Started timeutil.TimeStamp `xorm:"index"`
|
|
Stopped timeutil.TimeStamp `xorm:"index(stopped_log_expired)"`
|
|
|
|
RepoID int64 `xorm:"index"`
|
|
OwnerID int64 `xorm:"index"`
|
|
CommitSHA string `xorm:"index"`
|
|
IsForkPullRequest bool
|
|
|
|
Token string `xorm:"-"`
|
|
TokenHash string `xorm:"UNIQUE"` // sha256 of token
|
|
TokenSalt string
|
|
TokenLastEight string `xorm:"index token_last_eight"`
|
|
|
|
LogFilename string // file name of log
|
|
LogInStorage bool // read log from database or from storage
|
|
LogLength int64 // lines count
|
|
LogSize int64 // blob size
|
|
LogIndexes LogIndexes `xorm:"LONGBLOB"` // line number to offset
|
|
LogExpired bool `xorm:"index(stopped_log_expired)"` // files that are too old will be deleted
|
|
|
|
Created timeutil.TimeStamp `xorm:"created"`
|
|
Updated timeutil.TimeStamp `xorm:"updated index"`
|
|
}
|
|
|
|
// taskReportTimeout is how long a task may go without contact from its runner before the
|
|
// runner is assumed gone. Runners report state and stream logs every few seconds, both of
|
|
// which refresh ActionTask.Updated. Shorter than setting.Actions.ZombieTaskTimeout because
|
|
// it only decides whether the runner is reachable, not whether the task should be killed.
|
|
const taskReportTimeout = time.Minute
|
|
|
|
func init() {
|
|
db.RegisterModel(new(ActionTask))
|
|
}
|
|
|
|
func (task *ActionTask) IsStopped() bool {
|
|
return task.Stopped > 0
|
|
}
|
|
|
|
func (task *ActionTask) GetRunJobLink() string {
|
|
// Run.Repo can be nil when the repository was deleted while task/run rows remain
|
|
// (TaskList.LoadAttributes copies job.Repo into run.Repo, leaving it nil on a miss).
|
|
// Run.Link() already returns "" in that case, so guard here to avoid emitting a
|
|
// broken relative "/jobs/N" link from the Sprintf below.
|
|
if task.Job == nil || task.Job.Run == nil || task.Job.Run.Repo == nil {
|
|
return ""
|
|
}
|
|
return fmt.Sprintf("%s/jobs/%d", task.Job.Run.Link(), task.Job.ID)
|
|
}
|
|
|
|
func (task *ActionTask) GetCommitLink() string {
|
|
if task.Job == nil || task.Job.Run == nil || task.Job.Run.Repo == nil {
|
|
return ""
|
|
}
|
|
return task.Job.Run.Repo.CommitLink(task.CommitSHA)
|
|
}
|
|
|
|
func (task *ActionTask) GetRepoName() string {
|
|
if task.Job == nil || task.Job.Run == nil || task.Job.Run.Repo == nil {
|
|
return ""
|
|
}
|
|
return task.Job.Run.Repo.FullName()
|
|
}
|
|
|
|
func (task *ActionTask) GetRepoLink() string {
|
|
if task.Job == nil || task.Job.Run == nil || task.Job.Run.Repo == nil {
|
|
return ""
|
|
}
|
|
return task.Job.Run.Repo.Link()
|
|
}
|
|
|
|
func (task *ActionTask) LoadJob(ctx context.Context) error {
|
|
if task.Job == nil {
|
|
job, err := GetRunJobByRepoAndID(ctx, task.RepoID, task.JobID)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
task.Job = job
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// LoadAttributes load Job Steps if not loaded
|
|
func (task *ActionTask) LoadAttributes(ctx context.Context) error {
|
|
if err := task.LoadJob(ctx); err != nil {
|
|
return err
|
|
}
|
|
|
|
if err := task.Job.LoadAttributes(ctx); err != nil {
|
|
return err
|
|
}
|
|
|
|
if task.Steps == nil { // be careful, an empty slice (not nil) also means loaded
|
|
steps, err := GetTaskStepsByTaskID(ctx, task.ID)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
task.Steps = steps
|
|
}
|
|
|
|
return nil
|
|
}
|
|
|
|
func (task *ActionTask) GenerateAndFillToken() {
|
|
task.Token, task.TokenSalt, task.TokenHash, task.TokenLastEight = generateSaltedToken()
|
|
}
|
|
|
|
func GetTaskByID(ctx context.Context, id int64) (*ActionTask, error) {
|
|
var task ActionTask
|
|
has, err := db.GetEngine(ctx).Where("id=?", id).Get(&task)
|
|
if err != nil {
|
|
return nil, err
|
|
} else if !has {
|
|
return nil, fmt.Errorf("task with id %d: %w", id, util.ErrNotExist)
|
|
}
|
|
|
|
return &task, nil
|
|
}
|
|
|
|
// GetTasksMapByIDs returns the found tasks keyed by ID, silently omitting IDs that no longer exist.
|
|
func GetTasksMapByIDs(ctx context.Context, ids []int64) (map[int64]*ActionTask, error) {
|
|
tasks := make(map[int64]*ActionTask, len(ids))
|
|
if len(ids) == 0 {
|
|
return tasks, nil
|
|
}
|
|
return tasks, db.GetEngine(ctx).In("id", ids).Find(&tasks)
|
|
}
|
|
|
|
// GetTaskRunnerNames returns runner names keyed by task ID without loading task logs.
|
|
func GetTaskRunnerNames(ctx context.Context, taskIDs []int64) (map[int64]string, error) {
|
|
names := make(map[int64]string, len(taskIDs))
|
|
if len(taskIDs) == 0 {
|
|
return names, nil
|
|
}
|
|
var rows []struct {
|
|
ID int64
|
|
Name string
|
|
}
|
|
err := db.GetEngine(ctx).Table("action_task").
|
|
Join("INNER", "action_runner", "action_runner.id = action_task.runner_id").
|
|
In("action_task.id", taskIDs).
|
|
Select("action_task.id, action_runner.name").Find(&rows)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
for _, row := range rows {
|
|
names[row.ID] = row.Name
|
|
}
|
|
return names, nil
|
|
}
|
|
|
|
func GetRunningTaskByToken(ctx context.Context, token string) (*ActionTask, error) {
|
|
errNotExist := fmt.Errorf("task with token %q: %w", token, util.ErrNotExist)
|
|
if token == "" {
|
|
return nil, errNotExist
|
|
}
|
|
// A token is defined as being SHA1 sum these are 40 hexadecimal bytes long
|
|
if len(token) != 40 {
|
|
return nil, errNotExist
|
|
}
|
|
for _, x := range []byte(token) {
|
|
if x < '0' || (x > '9' && x < 'a') || x > 'f' {
|
|
return nil, errNotExist
|
|
}
|
|
}
|
|
|
|
cacheKey := "actions:" + token
|
|
lastEight := token[len(token)-8:]
|
|
if cached, _ := auth_model.TokenCache().Get(cacheKey); cached != nil {
|
|
task := &ActionTask{
|
|
TokenLastEight: lastEight,
|
|
}
|
|
// Re-get the task from the db in case it has been deleted in the intervening period
|
|
has, err := db.GetEngine(ctx).ID(cached.TokenID).Get(task)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
if has && util.CryptoConstTimeEqual(task.TokenHash, cached.TokenHash) {
|
|
return task, nil
|
|
}
|
|
auth_model.TokenCache().Remove(cacheKey)
|
|
}
|
|
|
|
var tasks []*ActionTask
|
|
// Cancelling tasks are still authenticating — post-run cleanup steps need API access (artifact uploads, cache saves, etc.) before the runner finalizes the task.
|
|
err := db.GetEngine(ctx).Where("token_last_eight = ? AND status IN (?, ?)", lastEight, StatusRunning, StatusCancelling).Find(&tasks)
|
|
if err != nil {
|
|
return nil, err
|
|
} else if len(tasks) == 0 {
|
|
return nil, errNotExist
|
|
}
|
|
|
|
for _, t := range tasks {
|
|
tempHash := auth_model.HashToken(token, t.TokenSalt)
|
|
if util.CryptoConstTimeEqual(t.TokenHash, tempHash) {
|
|
auth_model.TokenCache().Add(cacheKey, &auth_model.TokenCacheItem{TokenID: t.ID, TokenHash: t.TokenHash})
|
|
return t, nil
|
|
}
|
|
}
|
|
return nil, errNotExist
|
|
}
|
|
|
|
func makeTaskStepDisplayName(step *jobparser.Step, limit int) (name string) {
|
|
if step.Name != "" {
|
|
name = string(step.Name) // the step has an explicit name
|
|
} else {
|
|
// for unnamed step, its "String()" method tries to get a display name by its "name", "uses",
|
|
// "run" or "id" (last fallback), we add the "Run " prefix for unnamed steps for better display
|
|
// for multi-line "run" scripts, only use the first line to match GitHub's behavior
|
|
// https://github.com/actions/runner/blob/66800900843747f37591b077091dd2c8cf2c1796/src/Runner.Worker/Handlers/ScriptHandler.cs#L45-L58
|
|
runStr, _, _ := strings.Cut(strings.TrimSpace(string(step.Run)), "\n")
|
|
name = "Run " + util.IfZero(strings.TrimSpace(runStr), step.String())
|
|
}
|
|
return util.EllipsisDisplayString(name, limit) // database column has a length limit
|
|
}
|
|
|
|
// errJobAlreadyClaimed is a sentinel used inside claimJobForRunner to signal that
|
|
// another runner won the optimistic-lock race; it is never returned to callers.
|
|
var errJobAlreadyClaimed = errors.New("job already claimed by another runner")
|
|
|
|
// pickTaskBatchSize bounds how many waiting jobs each CreateTaskForRunner query loads,
|
|
// so a large backlog is not fetched into memory on every runner poll.
|
|
// It is a var only so tests can shrink it to exercise pagination cheaply.
|
|
var pickTaskBatchSize = 100
|
|
|
|
// CreateTaskForRunner finds a waiting job that matches the runner's labels and
|
|
// atomically claims it. It iterates through all matching jobs so that a
|
|
// concurrent claim by another runner (which would lose the optimistic lock on
|
|
// job #1) does not leave the remaining jobs permanently unassigned.
|
|
func CreateTaskForRunner(ctx context.Context, runner *ActionRunner) (*ActionTask, bool, error) {
|
|
if db.InTransaction(ctx) {
|
|
return nil, false, errors.New("CreateTaskForRunner must not be called within a database transaction")
|
|
}
|
|
e := db.GetEngine(ctx)
|
|
|
|
jobCond := builder.NewCond()
|
|
if runner.RepoID != 0 {
|
|
jobCond = builder.Eq{"repo_id": runner.RepoID}
|
|
} else if runner.OwnerID != 0 {
|
|
jobCond = builder.In("repo_id", builder.Select("`repository`.id").From("repository").
|
|
Join("INNER", "repo_unit", "`repository`.id = `repo_unit`.repo_id").
|
|
Where(builder.Eq{"`repository`.owner_id": runner.OwnerID, "`repo_unit`.type": unit.TypeActions}))
|
|
}
|
|
baseCond := builder.Eq{"task_id": 0, "status": StatusWaiting, "is_reusable_caller": false}.And(jobCond)
|
|
|
|
// TODO: a more efficient way to filter labels
|
|
log.Trace("runner labels: %v", runner.AgentLabels)
|
|
|
|
// Page through the waiting jobs oldest-first instead of loading the whole backlog into memory on every poll.
|
|
// Keyset pagination on (updated, id) is safe under concurrent claims:
|
|
// updated only moves forward, so the advancing cursor never skips a still-waiting job even as claimed jobs drop out.
|
|
var cursorUpdated timeutil.TimeStamp
|
|
var cursorID int64
|
|
for {
|
|
cond := baseCond
|
|
if cursorID > 0 {
|
|
cond = cond.And(builder.Or(
|
|
builder.Gt{"updated": cursorUpdated},
|
|
builder.And(builder.Eq{"updated": cursorUpdated}, builder.Gt{"id": cursorID}),
|
|
))
|
|
}
|
|
|
|
var jobs []*ActionRunJob
|
|
if err := e.Where(cond).Asc("updated", "id").Limit(pickTaskBatchSize).Find(&jobs); err != nil {
|
|
return nil, false, err
|
|
}
|
|
|
|
for _, v := range jobs {
|
|
if !runner.CanMatchLabels(v.RunsOn) {
|
|
continue
|
|
}
|
|
task, ok, err := claimJobForRunner(ctx, runner, v)
|
|
if err != nil {
|
|
return nil, false, err
|
|
}
|
|
if ok {
|
|
return task, true, nil
|
|
}
|
|
// Another runner claimed this job concurrently; try the next one.
|
|
}
|
|
|
|
// A short page means no waiting jobs remain beyond it.
|
|
if len(jobs) < pickTaskBatchSize {
|
|
return nil, false, nil
|
|
}
|
|
last := jobs[len(jobs)-1]
|
|
cursorUpdated, cursorID = last.Updated, last.ID
|
|
}
|
|
}
|
|
|
|
// claimJobForRunner attempts to atomically claim job for runner inside its own
|
|
// transaction. Returns (task, true, nil) on success, or (nil, false, nil) when
|
|
// another runner wins the optimistic-lock race (the caller should try the next
|
|
// candidate job).
|
|
func claimJobForRunner(ctx context.Context, runner *ActionRunner, job *ActionRunJob) (*ActionTask, bool, error) {
|
|
var resultTask *ActionTask
|
|
|
|
err := db.WithTx(ctx, func(ctx context.Context) error {
|
|
e := db.GetEngine(ctx)
|
|
|
|
if err := job.LoadAttributes(ctx); err != nil {
|
|
return err
|
|
}
|
|
|
|
now := timeutil.TimeStampNow()
|
|
job.Started = now
|
|
job.Status = StatusRunning
|
|
|
|
task := &ActionTask{
|
|
JobID: job.ID,
|
|
Attempt: job.Attempt,
|
|
RunnerID: runner.ID,
|
|
Started: now,
|
|
Status: StatusRunning,
|
|
RepoID: job.RepoID,
|
|
OwnerID: job.OwnerID,
|
|
CommitSHA: job.CommitSHA,
|
|
IsForkPullRequest: job.IsForkPullRequest,
|
|
}
|
|
task.GenerateAndFillToken()
|
|
|
|
workflowJob, err := job.ParseJob()
|
|
if err != nil {
|
|
return fmt.Errorf("load job %d: %w", job.ID, err)
|
|
}
|
|
|
|
if _, err := e.Insert(task); err != nil {
|
|
return err
|
|
}
|
|
|
|
task.LogFilename = logFileName(job.Run.Repo.FullName(), task.ID)
|
|
if err := UpdateTask(ctx, task, "log_filename"); err != nil {
|
|
return err
|
|
}
|
|
|
|
if len(workflowJob.Steps) > 0 {
|
|
steps := make([]*ActionTaskStep, len(workflowJob.Steps))
|
|
for i, v := range workflowJob.Steps {
|
|
steps[i] = &ActionTaskStep{
|
|
Name: makeTaskStepDisplayName(v, 255),
|
|
TaskID: task.ID,
|
|
Index: int64(i),
|
|
RepoID: task.RepoID,
|
|
Status: StatusWaiting,
|
|
}
|
|
}
|
|
if _, err := e.Insert(steps); err != nil {
|
|
return err
|
|
}
|
|
task.Steps = steps
|
|
}
|
|
|
|
job.TaskID = task.ID
|
|
n, err := UpdateRunJob(ctx, job, builder.And(builder.Eq{"task_id": 0}, builder.Eq{"status": StatusWaiting}))
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if n != 1 {
|
|
// Another runner claimed this job between our scan and this update;
|
|
// signal the outer loop to move on without treating this as an error.
|
|
return errJobAlreadyClaimed
|
|
}
|
|
|
|
task.Job = job
|
|
resultTask = task
|
|
return nil
|
|
})
|
|
|
|
if errors.Is(err, errJobAlreadyClaimed) {
|
|
return nil, false, nil
|
|
}
|
|
if err != nil {
|
|
return nil, false, err
|
|
}
|
|
return resultTask, true, nil
|
|
}
|
|
|
|
// ReleaseTaskForRunner reverts a freshly-claimed but undelivered task: it deletes
|
|
// the task together with its steps and returns the job to the waiting queue. It is
|
|
// used when assembling the runner response fails after the job was already claimed,
|
|
// so the job is not stranded in running state with no runner ever executing it.
|
|
func ReleaseTaskForRunner(ctx context.Context, task *ActionTask) error {
|
|
return db.WithTx(ctx, func(ctx context.Context) error {
|
|
e := db.GetEngine(ctx)
|
|
|
|
job, err := GetRunJobByRepoAndID(ctx, task.RepoID, task.JobID)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
job.Status = StatusWaiting
|
|
job.Started = 0
|
|
job.TaskID = 0
|
|
// Guard on task_id and status so we only release while the job still
|
|
// references this task and has not progressed past running.
|
|
n, err := UpdateRunJob(ctx, job, builder.Eq{"task_id": task.ID, "status": StatusRunning}, "status", "started", "task_id")
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if n != 1 {
|
|
return fmt.Errorf("release task %d: job %d no longer references it", task.ID, task.JobID)
|
|
}
|
|
|
|
if _, err := e.Delete(&ActionTaskStep{TaskID: task.ID}); err != nil {
|
|
return err
|
|
}
|
|
if _, err := e.ID(task.ID).Delete(&ActionTask{}); err != nil {
|
|
return err
|
|
}
|
|
return nil
|
|
})
|
|
}
|
|
|
|
func UpdateTask(ctx context.Context, task *ActionTask, cols ...string) error {
|
|
sess := db.GetEngine(ctx).ID(task.ID)
|
|
if len(cols) > 0 {
|
|
sess.Cols(cols...)
|
|
}
|
|
_, err := sess.Update(task)
|
|
|
|
// Automatically delete the ephemeral runner if the task is done
|
|
if err == nil && task.Status.IsDone() && util.SliceContainsString(cols, "status") {
|
|
return DeleteEphemeralRunner(ctx, task.RunnerID)
|
|
}
|
|
return err
|
|
}
|
|
|
|
func getRunIDByTaskID(ctx context.Context, taskID int64) (runID int64, _ error) {
|
|
if has, err := db.GetEngine(ctx).Cols("action_run_job.run_id").
|
|
Table("action_task").
|
|
Join("INNER", "action_run_job", "action_run_job.id = action_task.job_id").
|
|
Where(builder.Eq{"action_task.id": taskID}).Get(&runID); err != nil {
|
|
return runID, err
|
|
} else if !has {
|
|
return runID, util.ErrNotExist
|
|
}
|
|
return runID, nil
|
|
}
|
|
|
|
// UpdateTaskByState updates the task by the state.
|
|
// It will always update the task if the state is not final, even there is no change.
|
|
// So it will update ActionTask.Updated to avoid the task being judged as a zombie task.
|
|
func UpdateTaskByState(ctx context.Context, runnerID int64, state *runnerv1.TaskState) (*ActionTask, error) {
|
|
stepStates := map[int64]*runnerv1.StepState{}
|
|
for _, v := range state.Steps {
|
|
stepStates[v.Id] = v
|
|
}
|
|
|
|
// Only one request can update the task because the final state needs to be calculated with all job states.
|
|
// Otherwise, concurrent requests with transaction will make the SQL read stale job state and result in wrong final state.
|
|
taskID := state.Id
|
|
runID, err := getRunIDByTaskID(ctx, taskID)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
task := &ActionTask{}
|
|
applyState := func(ctx context.Context) error {
|
|
if has, err := db.GetEngine(ctx).ID(taskID).Get(task); err != nil {
|
|
return err
|
|
} else if !has {
|
|
return util.ErrNotExist
|
|
} else if runnerID != task.RunnerID {
|
|
return errors.New("invalid runner for task")
|
|
}
|
|
|
|
if task.Status.IsDone() {
|
|
// the state is final, do nothing
|
|
return nil
|
|
}
|
|
|
|
now := timeutil.TimeStampNow()
|
|
// state.Result is not unspecified means the task is finished
|
|
if state.Result != runnerv1.Result_RESULT_UNSPECIFIED {
|
|
if task.Status == StatusCancelling {
|
|
// The runner may report SUCCESS/FAILURE for the cleanup phase; preserve user intent.
|
|
task.Status = StatusCancelled
|
|
} else {
|
|
task.Status = StatusFromResult(state.Result)
|
|
}
|
|
task.Stopped = now
|
|
if err := UpdateTask(ctx, task, "status", "stopped"); err != nil {
|
|
return err
|
|
}
|
|
if _, err := UpdateRunJob(ctx, &ActionRunJob{
|
|
ID: task.JobID,
|
|
RepoID: task.RepoID,
|
|
Status: task.Status,
|
|
Stopped: task.Stopped,
|
|
}, nil, "status", "stopped"); err != nil {
|
|
return err
|
|
}
|
|
} else {
|
|
// Force update ActionTask.Updated to avoid the task being judged as a zombie task
|
|
task.Updated = now
|
|
if err := UpdateTask(ctx, task, "updated"); err != nil {
|
|
return err
|
|
}
|
|
}
|
|
|
|
if err := task.LoadAttributes(ctx); err != nil {
|
|
return err
|
|
}
|
|
|
|
for _, step := range task.Steps {
|
|
var result runnerv1.Result
|
|
if v, ok := stepStates[step.Index]; ok {
|
|
result = v.Result
|
|
step.LogIndex = v.LogIndex
|
|
step.LogLength = v.LogLength
|
|
if step.Started == 0 && v.StartedAt != nil {
|
|
step.Started = now
|
|
}
|
|
}
|
|
if result != runnerv1.Result_RESULT_UNSPECIFIED {
|
|
step.Status = StatusFromResult(result)
|
|
step.Stopped = util.IfZero(step.Stopped, now)
|
|
} else if step.Started != 0 {
|
|
step.Status = StatusRunning
|
|
}
|
|
if _, err := db.GetEngine(ctx).ID(step.ID).Update(step); err != nil {
|
|
return err
|
|
}
|
|
}
|
|
return nil
|
|
}
|
|
err = globallock.LockAndDo(ctx, fmt.Sprintf("UpdateTaskByState-run-%d", runID), func(ctx context.Context) error {
|
|
// A half-written report leaves the task done with a running job, which no retry repairs.
|
|
return db.WithTx(ctx, applyState)
|
|
})
|
|
return task, err
|
|
}
|
|
|
|
func StopTask(ctx context.Context, taskID int64, status Status) error {
|
|
if !status.IsDone() && status != StatusCancelling {
|
|
return fmt.Errorf("cannot stop task with status %v", status)
|
|
}
|
|
e := db.GetEngine(ctx)
|
|
|
|
task := &ActionTask{}
|
|
if has, err := e.ID(taskID).Get(task); err != nil {
|
|
return err
|
|
} else if !has {
|
|
return util.ErrNotExist
|
|
}
|
|
if task.Status.IsDone() {
|
|
return nil
|
|
}
|
|
|
|
now := timeutil.TimeStampNow()
|
|
if status == StatusCancelling {
|
|
runner, err := GetRunnerByID(ctx, task.RunnerID)
|
|
if err != nil {
|
|
if !errors.Is(err, util.ErrNotExist) {
|
|
return err
|
|
}
|
|
status = StatusCancelled
|
|
} else if !runner.HasCancellingSupport {
|
|
status = StatusCancelled
|
|
} else if task.Updated.AddDuration(taskReportTimeout) < now {
|
|
// A runner that stopped reporting will never acknowledge the cancellation either,
|
|
// so skip the handshake instead of waiting for the zombie task cleanup.
|
|
status = StatusCancelled
|
|
}
|
|
}
|
|
|
|
if status == StatusCancelling {
|
|
task.Status = StatusCancelling
|
|
|
|
if _, err := UpdateRunJob(ctx, &ActionRunJob{
|
|
ID: task.JobID,
|
|
RepoID: task.RepoID,
|
|
Status: StatusCancelling,
|
|
}, nil, "status"); err != nil {
|
|
return err
|
|
}
|
|
|
|
// NoAutoTime keeps "updated" at the runner's last contact: re-cancelling an already
|
|
// cancelling task must not defer the timeout above or the zombie task cleanup.
|
|
_, err := e.ID(task.ID).Cols("status").NoAutoTime().Update(task)
|
|
return err
|
|
}
|
|
|
|
task.Status = status
|
|
task.Stopped = now
|
|
if _, err := UpdateRunJob(ctx, &ActionRunJob{
|
|
ID: task.JobID,
|
|
RepoID: task.RepoID,
|
|
Status: task.Status,
|
|
Stopped: task.Stopped,
|
|
}, nil); err != nil {
|
|
return err
|
|
}
|
|
|
|
if err := UpdateTask(ctx, task, "status", "stopped"); err != nil {
|
|
return err
|
|
}
|
|
|
|
if err := task.LoadAttributes(ctx); err != nil {
|
|
return err
|
|
}
|
|
|
|
for _, step := range task.Steps {
|
|
if !step.Status.IsDone() {
|
|
step.Status = status
|
|
if step.Started == 0 {
|
|
step.Started = now
|
|
}
|
|
step.Stopped = now
|
|
}
|
|
if _, err := e.ID(step.ID).Update(step); err != nil {
|
|
return err
|
|
}
|
|
}
|
|
|
|
return nil
|
|
}
|
|
|
|
func FindOldTasksToExpire(ctx context.Context, olderThan timeutil.TimeStamp, limit int) ([]*ActionTask, error) {
|
|
e := db.GetEngine(ctx)
|
|
|
|
tasks := make([]*ActionTask, 0, limit)
|
|
// Check "stopped > 0" to avoid deleting tasks that are still running
|
|
return tasks, e.Where("stopped > 0 AND stopped < ? AND log_expired = ?", olderThan, false).
|
|
Limit(limit).
|
|
Find(&tasks)
|
|
}
|
|
|
|
func logFileName(repoFullName string, taskID int64) string {
|
|
ret := fmt.Sprintf("%s/%02x/%d.log", repoFullName, taskID%256, taskID)
|
|
|
|
if setting.Actions.LogCompression.IsZstd() {
|
|
ret += ".zst"
|
|
}
|
|
|
|
return ret
|
|
}
|