package store import ( "context" "database/sql" "errors" "strings" "time" ) // Job states, where a lost runner requeues once, because an infinite requeue is the classic failure. const ( JobQueued = "queued" JobRunning = "running" JobDone = "done" JobLost = "lost" ) type Job struct { ID int64 Repo string Ref string SHA string Command string Image string // Labels is the machine this job asked for, empty meaning any. Chapter 15A. Labels []string // Name is the workflow job and its matrix combination, so two builds of one commit are told apart. Name string State string RunnerID int64 Attempts int Created time.Time Started time.Time Log string } // QueueJob adds a build to a repository's queue. func (db *DB) QueueJob(ctx context.Context, repo, ref, sha, command, image string) (int64, error) { return db.QueueJobFor(ctx, repo, ref, sha, command, image, nil, "") } // QueueJobFor queues a job only a machine with these labels may take, under a name. Chapter 15A. func (db *DB) QueueJobFor(ctx context.Context, repo, ref, sha, command, image string, labels []string, name string) (int64, error) { var id int64 err := db.QueryRowContext(ctx, `INSERT INTO jobs (repo, ref, sha, command, image, labels, name, state, created_at) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?) RETURNING id`, repo, ref, sha, command, image, strings.Join(labels, ","), name, JobQueued, now().Unix()).Scan(&id) return id, err } // TakeJob hands over the oldest queued job, first in and first out. Chapter 15. func (db *DB) TakeJob(ctx context.Context, repo string, runner Runner) (*Job, error) { tx, err := db.Begin(ctx) if err != nil { return nil, err } defer tx.Rollback() // The first job this machine can run is taken, so one waiting for another machine blocks nothing. rows, err := tx.QueryContext(ctx, `SELECT id, repo, ref, sha, command, image, labels, name, attempts, created_at FROM jobs WHERE repo = ? AND state = ? ORDER BY id`, repo, JobQueued) if err != nil { return nil, err } var j Job var created int64 found := false for rows.Next() { var cand Job var labels string var at int64 if err := rows.Scan(&cand.ID, &cand.Repo, &cand.Ref, &cand.SHA, &cand.Command, &cand.Image, &labels, &cand.Name, &cand.Attempts, &at); err != nil { rows.Close() return nil, err } cand.Labels = splitLabels(labels) if !runner.Satisfies(cand.Labels) { continue } j, created, found = cand, at, true break } rows.Close() if err := rows.Err(); err != nil { return nil, err } if !found { return nil, nil } j.Created = time.Unix(created, 0) j.State = JobRunning j.RunnerID = runner.ID j.Started = now() // Still queued is part of the write, so two machines reaching for one job leave one holding it. res, err := tx.ExecContext(ctx, `UPDATE jobs SET state = ?, runner_id = ?, started_at = ?, attempts = attempts + 1 WHERE id = ? AND state = ?`, JobRunning, runner.ID, j.Started.Unix(), j.ID, JobQueued) if err != nil { return nil, err } took, err := res.RowsAffected() if err != nil { return nil, err } if took == 0 { // Somebody else took it between the read and the write, and the runner asks again on its tick. return nil, nil } if err := tx.Commit(); err != nil { return nil, err } return &j, nil } // Job reads one job. func (db *DB) Job(ctx context.Context, id int64) (*Job, error) { var j Job var runner sql.NullInt64 var created int64 var started sql.NullInt64 err := db.QueryRowContext(ctx, `SELECT id, repo, ref, sha, command, image, name, state, runner_id, attempts, created_at, started_at, log FROM jobs WHERE id = ?`, id). Scan(&j.ID, &j.Repo, &j.Ref, &j.SHA, &j.Command, &j.Image, &j.Name, &j.State, &runner, &j.Attempts, &created, &started, &j.Log) if errors.Is(err, sql.ErrNoRows) { return nil, ErrNotFound } if err != nil { return nil, err } j.RunnerID = runner.Int64 j.Created = time.Unix(created, 0) if started.Valid { j.Started = time.Unix(started.Int64, 0) } return &j, nil } // AppendLog adds a chunk of build output to a running job. func (db *DB) AppendLog(ctx context.Context, id int64, runnerID int64, chunk string) error { res, err := db.ExecContext(ctx, `UPDATE jobs SET log = log || ? WHERE id = ? AND runner_id = ? AND state = ?`, chunk, id, runnerID, JobRunning) if err != nil { return err } if n, _ := res.RowsAffected(); n == 0 { return ErrNotFound } return nil } // FinishJob closes a job out and returns it, so the caller can write the note. func (db *DB) FinishJob(ctx context.Context, id, runnerID int64) (*Job, error) { res, err := db.ExecContext(ctx, `UPDATE jobs SET state = ? WHERE id = ? AND runner_id = ? AND state = ?`, JobDone, id, runnerID, JobRunning) if err != nil { return nil, err } if n, _ := res.RowsAffected(); n == 0 { return nil, ErrNotFound } return db.Job(ctx, id) } // RequeueLostJobs retries a silent runner's job once, then marks it lost rather than cycling. func (db *DB) RequeueLostJobs(ctx context.Context, olderThan time.Duration) (int, error) { cutoff := now().Add(-olderThan).Unix() res, err := db.ExecContext(ctx, `UPDATE jobs SET state = ?, runner_id = NULL WHERE state = ? AND started_at < ? AND attempts < 2`, JobQueued, JobRunning, cutoff) if err != nil { return 0, err } requeued, _ := res.RowsAffected() if _, err := db.ExecContext(ctx, `UPDATE jobs SET state = ? WHERE state = ? AND started_at < ? AND attempts >= 2`, JobLost, JobRunning, cutoff); err != nil { return int(requeued), err } return int(requeued), nil } // QueuedJobs lists what is waiting for a machine, which the runs page and the tests both ask for. func (db *DB) QueuedJobs(ctx context.Context, repo string) ([]Job, error) { rows, err := db.QueryContext(ctx, `SELECT id, repo, ref, sha, command, image, labels, name, attempts, created_at FROM jobs WHERE repo = ? AND state = ? ORDER BY id`, repo, JobQueued) if err != nil { return nil, err } defer rows.Close() var out []Job for rows.Next() { var j Job var labels string var created int64 if err := rows.Scan(&j.ID, &j.Repo, &j.Ref, &j.SHA, &j.Command, &j.Image, &labels, &j.Name, &j.Attempts, &created); err != nil { return nil, err } j.Labels = splitLabels(labels) j.State = JobQueued j.Created = time.Unix(created, 0) out = append(out, j) } return out, rows.Err() }