gitea源码

task.go 13KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527
  1. // Copyright 2022 The Gitea Authors. All rights reserved.
  2. // SPDX-License-Identifier: MIT
  3. package actions
  4. import (
  5. "context"
  6. "crypto/subtle"
  7. "errors"
  8. "fmt"
  9. "time"
  10. auth_model "code.gitea.io/gitea/models/auth"
  11. "code.gitea.io/gitea/models/db"
  12. "code.gitea.io/gitea/models/unit"
  13. "code.gitea.io/gitea/modules/container"
  14. "code.gitea.io/gitea/modules/log"
  15. "code.gitea.io/gitea/modules/setting"
  16. "code.gitea.io/gitea/modules/timeutil"
  17. "code.gitea.io/gitea/modules/util"
  18. runnerv1 "code.gitea.io/actions-proto-go/runner/v1"
  19. lru "github.com/hashicorp/golang-lru/v2"
  20. "github.com/nektos/act/pkg/jobparser"
  21. "google.golang.org/protobuf/types/known/timestamppb"
  22. "xorm.io/builder"
  23. )
  24. // ActionTask represents a distribution of job
  25. type ActionTask struct {
  26. ID int64
  27. JobID int64
  28. Job *ActionRunJob `xorm:"-"`
  29. Steps []*ActionTaskStep `xorm:"-"`
  30. Attempt int64
  31. RunnerID int64 `xorm:"index"`
  32. Status Status `xorm:"index"`
  33. Started timeutil.TimeStamp `xorm:"index"`
  34. Stopped timeutil.TimeStamp `xorm:"index(stopped_log_expired)"`
  35. RepoID int64 `xorm:"index"`
  36. OwnerID int64 `xorm:"index"`
  37. CommitSHA string `xorm:"index"`
  38. IsForkPullRequest bool
  39. Token string `xorm:"-"`
  40. TokenHash string `xorm:"UNIQUE"` // sha256 of token
  41. TokenSalt string
  42. TokenLastEight string `xorm:"index token_last_eight"`
  43. LogFilename string // file name of log
  44. LogInStorage bool // read log from database or from storage
  45. LogLength int64 // lines count
  46. LogSize int64 // blob size
  47. LogIndexes LogIndexes `xorm:"LONGBLOB"` // line number to offset
  48. LogExpired bool `xorm:"index(stopped_log_expired)"` // files that are too old will be deleted
  49. Created timeutil.TimeStamp `xorm:"created"`
  50. Updated timeutil.TimeStamp `xorm:"updated index"`
  51. }
  52. var successfulTokenTaskCache *lru.Cache[string, any]
  53. func init() {
  54. db.RegisterModel(new(ActionTask), func() error {
  55. if setting.SuccessfulTokensCacheSize > 0 {
  56. var err error
  57. successfulTokenTaskCache, err = lru.New[string, any](setting.SuccessfulTokensCacheSize)
  58. if err != nil {
  59. return fmt.Errorf("unable to allocate Task cache: %v", err)
  60. }
  61. } else {
  62. successfulTokenTaskCache = nil
  63. }
  64. return nil
  65. })
  66. }
  67. func (task *ActionTask) Duration() time.Duration {
  68. return calculateDuration(task.Started, task.Stopped, task.Status)
  69. }
  70. func (task *ActionTask) IsStopped() bool {
  71. return task.Stopped > 0
  72. }
  73. func (task *ActionTask) GetRunLink() string {
  74. if task.Job == nil || task.Job.Run == nil {
  75. return ""
  76. }
  77. return task.Job.Run.Link()
  78. }
  79. func (task *ActionTask) GetCommitLink() string {
  80. if task.Job == nil || task.Job.Run == nil || task.Job.Run.Repo == nil {
  81. return ""
  82. }
  83. return task.Job.Run.Repo.CommitLink(task.CommitSHA)
  84. }
  85. func (task *ActionTask) GetRepoName() string {
  86. if task.Job == nil || task.Job.Run == nil || task.Job.Run.Repo == nil {
  87. return ""
  88. }
  89. return task.Job.Run.Repo.FullName()
  90. }
  91. func (task *ActionTask) GetRepoLink() string {
  92. if task.Job == nil || task.Job.Run == nil || task.Job.Run.Repo == nil {
  93. return ""
  94. }
  95. return task.Job.Run.Repo.Link()
  96. }
  97. func (task *ActionTask) LoadJob(ctx context.Context) error {
  98. if task.Job == nil {
  99. job, err := GetRunJobByID(ctx, task.JobID)
  100. if err != nil {
  101. return err
  102. }
  103. task.Job = job
  104. }
  105. return nil
  106. }
  107. // LoadAttributes load Job Steps if not loaded
  108. func (task *ActionTask) LoadAttributes(ctx context.Context) error {
  109. if task == nil {
  110. return nil
  111. }
  112. if err := task.LoadJob(ctx); err != nil {
  113. return err
  114. }
  115. if err := task.Job.LoadAttributes(ctx); err != nil {
  116. return err
  117. }
  118. if task.Steps == nil { // be careful, an empty slice (not nil) also means loaded
  119. steps, err := GetTaskStepsByTaskID(ctx, task.ID)
  120. if err != nil {
  121. return err
  122. }
  123. task.Steps = steps
  124. }
  125. return nil
  126. }
  127. func (task *ActionTask) GenerateToken() (err error) {
  128. task.Token, task.TokenSalt, task.TokenHash, task.TokenLastEight, err = generateSaltedToken()
  129. return err
  130. }
  131. func GetTaskByID(ctx context.Context, id int64) (*ActionTask, error) {
  132. var task ActionTask
  133. has, err := db.GetEngine(ctx).Where("id=?", id).Get(&task)
  134. if err != nil {
  135. return nil, err
  136. } else if !has {
  137. return nil, fmt.Errorf("task with id %d: %w", id, util.ErrNotExist)
  138. }
  139. return &task, nil
  140. }
  141. func GetRunningTaskByToken(ctx context.Context, token string) (*ActionTask, error) {
  142. errNotExist := fmt.Errorf("task with token %q: %w", token, util.ErrNotExist)
  143. if token == "" {
  144. return nil, errNotExist
  145. }
  146. // A token is defined as being SHA1 sum these are 40 hexadecimal bytes long
  147. if len(token) != 40 {
  148. return nil, errNotExist
  149. }
  150. for _, x := range []byte(token) {
  151. if x < '0' || (x > '9' && x < 'a') || x > 'f' {
  152. return nil, errNotExist
  153. }
  154. }
  155. lastEight := token[len(token)-8:]
  156. if id := getTaskIDFromCache(token); id > 0 {
  157. task := &ActionTask{
  158. TokenLastEight: lastEight,
  159. }
  160. // Re-get the task from the db in case it has been deleted in the intervening period
  161. has, err := db.GetEngine(ctx).ID(id).Get(task)
  162. if err != nil {
  163. return nil, err
  164. }
  165. if has {
  166. return task, nil
  167. }
  168. successfulTokenTaskCache.Remove(token)
  169. }
  170. var tasks []*ActionTask
  171. err := db.GetEngine(ctx).Where("token_last_eight = ? AND status = ?", lastEight, StatusRunning).Find(&tasks)
  172. if err != nil {
  173. return nil, err
  174. } else if len(tasks) == 0 {
  175. return nil, errNotExist
  176. }
  177. for _, t := range tasks {
  178. tempHash := auth_model.HashToken(token, t.TokenSalt)
  179. if subtle.ConstantTimeCompare([]byte(t.TokenHash), []byte(tempHash)) == 1 {
  180. if successfulTokenTaskCache != nil {
  181. successfulTokenTaskCache.Add(token, t.ID)
  182. }
  183. return t, nil
  184. }
  185. }
  186. return nil, errNotExist
  187. }
  188. func CreateTaskForRunner(ctx context.Context, runner *ActionRunner) (*ActionTask, bool, error) {
  189. ctx, committer, err := db.TxContext(ctx)
  190. if err != nil {
  191. return nil, false, err
  192. }
  193. defer committer.Close()
  194. e := db.GetEngine(ctx)
  195. jobCond := builder.NewCond()
  196. if runner.RepoID != 0 {
  197. jobCond = builder.Eq{"repo_id": runner.RepoID}
  198. } else if runner.OwnerID != 0 {
  199. jobCond = builder.In("repo_id", builder.Select("`repository`.id").From("repository").
  200. Join("INNER", "repo_unit", "`repository`.id = `repo_unit`.repo_id").
  201. Where(builder.Eq{"`repository`.owner_id": runner.OwnerID, "`repo_unit`.type": unit.TypeActions}))
  202. }
  203. if jobCond.IsValid() {
  204. jobCond = builder.In("run_id", builder.Select("id").From("action_run").Where(jobCond))
  205. }
  206. var jobs []*ActionRunJob
  207. if err := e.Where("task_id=? AND status=?", 0, StatusWaiting).And(jobCond).Asc("updated", "id").Find(&jobs); err != nil {
  208. return nil, false, err
  209. }
  210. // TODO: a more efficient way to filter labels
  211. var job *ActionRunJob
  212. log.Trace("runner labels: %v", runner.AgentLabels)
  213. for _, v := range jobs {
  214. if isSubset(runner.AgentLabels, v.RunsOn) {
  215. job = v
  216. break
  217. }
  218. }
  219. if job == nil {
  220. return nil, false, nil
  221. }
  222. if err := job.LoadAttributes(ctx); err != nil {
  223. return nil, false, err
  224. }
  225. now := timeutil.TimeStampNow()
  226. job.Attempt++
  227. job.Started = now
  228. job.Status = StatusRunning
  229. task := &ActionTask{
  230. JobID: job.ID,
  231. Attempt: job.Attempt,
  232. RunnerID: runner.ID,
  233. Started: now,
  234. Status: StatusRunning,
  235. RepoID: job.RepoID,
  236. OwnerID: job.OwnerID,
  237. CommitSHA: job.CommitSHA,
  238. IsForkPullRequest: job.IsForkPullRequest,
  239. }
  240. if err := task.GenerateToken(); err != nil {
  241. return nil, false, err
  242. }
  243. parsedWorkflows, err := jobparser.Parse(job.WorkflowPayload)
  244. if err != nil {
  245. return nil, false, fmt.Errorf("parse workflow of job %d: %w", job.ID, err)
  246. } else if len(parsedWorkflows) != 1 {
  247. return nil, false, fmt.Errorf("workflow of job %d: not single workflow", job.ID)
  248. }
  249. _, workflowJob := parsedWorkflows[0].Job()
  250. if _, err := e.Insert(task); err != nil {
  251. return nil, false, err
  252. }
  253. task.LogFilename = logFileName(job.Run.Repo.FullName(), task.ID)
  254. if err := UpdateTask(ctx, task, "log_filename"); err != nil {
  255. return nil, false, err
  256. }
  257. if len(workflowJob.Steps) > 0 {
  258. steps := make([]*ActionTaskStep, len(workflowJob.Steps))
  259. for i, v := range workflowJob.Steps {
  260. name := util.EllipsisDisplayString(v.String(), 255)
  261. steps[i] = &ActionTaskStep{
  262. Name: name,
  263. TaskID: task.ID,
  264. Index: int64(i),
  265. RepoID: task.RepoID,
  266. Status: StatusWaiting,
  267. }
  268. }
  269. if _, err := e.Insert(steps); err != nil {
  270. return nil, false, err
  271. }
  272. task.Steps = steps
  273. }
  274. job.TaskID = task.ID
  275. if n, err := UpdateRunJob(ctx, job, builder.Eq{"task_id": 0}); err != nil {
  276. return nil, false, err
  277. } else if n != 1 {
  278. return nil, false, nil
  279. }
  280. task.Job = job
  281. if err := committer.Commit(); err != nil {
  282. return nil, false, err
  283. }
  284. return task, true, nil
  285. }
  286. func UpdateTask(ctx context.Context, task *ActionTask, cols ...string) error {
  287. sess := db.GetEngine(ctx).ID(task.ID)
  288. if len(cols) > 0 {
  289. sess.Cols(cols...)
  290. }
  291. _, err := sess.Update(task)
  292. // Automatically delete the ephemeral runner if the task is done
  293. if err == nil && task.Status.IsDone() && util.SliceContainsString(cols, "status") {
  294. return DeleteEphemeralRunner(ctx, task.RunnerID)
  295. }
  296. return err
  297. }
  298. // UpdateTaskByState updates the task by the state.
  299. // It will always update the task if the state is not final, even there is no change.
  300. // So it will update ActionTask.Updated to avoid the task being judged as a zombie task.
  301. func UpdateTaskByState(ctx context.Context, runnerID int64, state *runnerv1.TaskState) (*ActionTask, error) {
  302. stepStates := map[int64]*runnerv1.StepState{}
  303. for _, v := range state.Steps {
  304. stepStates[v.Id] = v
  305. }
  306. return db.WithTx2(ctx, func(ctx context.Context) (*ActionTask, error) {
  307. e := db.GetEngine(ctx)
  308. task := &ActionTask{}
  309. if has, err := e.ID(state.Id).Get(task); err != nil {
  310. return nil, err
  311. } else if !has {
  312. return nil, util.ErrNotExist
  313. } else if runnerID != task.RunnerID {
  314. return nil, errors.New("invalid runner for task")
  315. }
  316. if task.Status.IsDone() {
  317. // the state is final, do nothing
  318. return task, nil
  319. }
  320. // state.Result is not unspecified means the task is finished
  321. if state.Result != runnerv1.Result_RESULT_UNSPECIFIED {
  322. task.Status = Status(state.Result)
  323. task.Stopped = timeutil.TimeStamp(state.StoppedAt.AsTime().Unix())
  324. if err := UpdateTask(ctx, task, "status", "stopped"); err != nil {
  325. return nil, err
  326. }
  327. if _, err := UpdateRunJob(ctx, &ActionRunJob{
  328. ID: task.JobID,
  329. Status: task.Status,
  330. Stopped: task.Stopped,
  331. }, nil); err != nil {
  332. return nil, err
  333. }
  334. } else {
  335. // Force update ActionTask.Updated to avoid the task being judged as a zombie task
  336. task.Updated = timeutil.TimeStampNow()
  337. if err := UpdateTask(ctx, task, "updated"); err != nil {
  338. return nil, err
  339. }
  340. }
  341. if err := task.LoadAttributes(ctx); err != nil {
  342. return nil, err
  343. }
  344. for _, step := range task.Steps {
  345. var result runnerv1.Result
  346. if v, ok := stepStates[step.Index]; ok {
  347. result = v.Result
  348. step.LogIndex = v.LogIndex
  349. step.LogLength = v.LogLength
  350. step.Started = convertTimestamp(v.StartedAt)
  351. step.Stopped = convertTimestamp(v.StoppedAt)
  352. }
  353. if result != runnerv1.Result_RESULT_UNSPECIFIED {
  354. step.Status = Status(result)
  355. } else if step.Started != 0 {
  356. step.Status = StatusRunning
  357. }
  358. if _, err := e.ID(step.ID).Update(step); err != nil {
  359. return nil, err
  360. }
  361. }
  362. return task, nil
  363. })
  364. }
  365. func StopTask(ctx context.Context, taskID int64, status Status) error {
  366. if !status.IsDone() {
  367. return fmt.Errorf("cannot stop task with status %v", status)
  368. }
  369. e := db.GetEngine(ctx)
  370. task := &ActionTask{}
  371. if has, err := e.ID(taskID).Get(task); err != nil {
  372. return err
  373. } else if !has {
  374. return util.ErrNotExist
  375. }
  376. if task.Status.IsDone() {
  377. return nil
  378. }
  379. now := timeutil.TimeStampNow()
  380. task.Status = status
  381. task.Stopped = now
  382. if _, err := UpdateRunJob(ctx, &ActionRunJob{
  383. ID: task.JobID,
  384. Status: task.Status,
  385. Stopped: task.Stopped,
  386. }, nil); err != nil {
  387. return err
  388. }
  389. if err := UpdateTask(ctx, task, "status", "stopped"); err != nil {
  390. return err
  391. }
  392. if err := task.LoadAttributes(ctx); err != nil {
  393. return err
  394. }
  395. for _, step := range task.Steps {
  396. if !step.Status.IsDone() {
  397. step.Status = status
  398. if step.Started == 0 {
  399. step.Started = now
  400. }
  401. step.Stopped = now
  402. }
  403. if _, err := e.ID(step.ID).Update(step); err != nil {
  404. return err
  405. }
  406. }
  407. return nil
  408. }
  409. func FindOldTasksToExpire(ctx context.Context, olderThan timeutil.TimeStamp, limit int) ([]*ActionTask, error) {
  410. e := db.GetEngine(ctx)
  411. tasks := make([]*ActionTask, 0, limit)
  412. // Check "stopped > 0" to avoid deleting tasks that are still running
  413. return tasks, e.Where("stopped > 0 AND stopped < ? AND log_expired = ?", olderThan, false).
  414. Limit(limit).
  415. Find(&tasks)
  416. }
  417. func isSubset(set, subset []string) bool {
  418. m := make(container.Set[string], len(set))
  419. for _, v := range set {
  420. m.Add(v)
  421. }
  422. for _, v := range subset {
  423. if !m.Contains(v) {
  424. return false
  425. }
  426. }
  427. return true
  428. }
  429. func convertTimestamp(timestamp *timestamppb.Timestamp) timeutil.TimeStamp {
  430. if timestamp.GetSeconds() == 0 && timestamp.GetNanos() == 0 {
  431. return timeutil.TimeStamp(0)
  432. }
  433. return timeutil.TimeStamp(timestamp.AsTime().Unix())
  434. }
  435. func logFileName(repoFullName string, taskID int64) string {
  436. ret := fmt.Sprintf("%s/%02x/%d.log", repoFullName, taskID%256, taskID)
  437. if setting.Actions.LogCompression.IsZstd() {
  438. ret += ".zst"
  439. }
  440. return ret
  441. }
  442. func getTaskIDFromCache(token string) int64 {
  443. if successfulTokenTaskCache == nil {
  444. return 0
  445. }
  446. tInterface, ok := successfulTokenTaskCache.Get(token)
  447. if !ok {
  448. return 0
  449. }
  450. t, ok := tInterface.(int64)
  451. if !ok {
  452. return 0
  453. }
  454. return t
  455. }