You can not select more than 25 topics Topics must start with a letter or number, can include dashes ('-') and can be up to 35 characters long.

task.go 12KB

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