File view with blame information shown in the left gutter beside each line.

barerepo / server / internal/store/jobs.go
220 lines · 6.3kb · 133728eaa05486504991b230f9d4c7e986b7defc
log files threads runs releases config jump to file t
133728e barerepo 1mo
1
package store
133728e barerepo 1mo
2
133728e barerepo 1mo
3
import (
133728e barerepo 1mo
4
"context"
133728e barerepo 1mo
5
"database/sql"
133728e barerepo 1mo
6
"errors"
133728e barerepo 1mo
7
"strings"
133728e barerepo 1mo
8
"time"
133728e barerepo 1mo
9
)
133728e barerepo 1mo
10
133728e barerepo 1mo
11
// Job states, where a lost runner requeues once, because an infinite requeue is the classic failure.
133728e barerepo 1mo
12
const (
133728e barerepo 1mo
13
JobQueued = "queued"
133728e barerepo 1mo
14
JobRunning = "running"
133728e barerepo 1mo
15
JobDone = "done"
133728e barerepo 1mo
16
JobLost = "lost"
133728e barerepo 1mo
17
)
133728e barerepo 1mo
18
133728e barerepo 1mo
19
type Job struct {
133728e barerepo 1mo
20
ID int64
133728e barerepo 1mo
21
Repo string
133728e barerepo 1mo
22
Ref string
133728e barerepo 1mo
23
SHA string
133728e barerepo 1mo
24
Command string
133728e barerepo 1mo
25
Image string
133728e barerepo 1mo
26
// Labels is the machine this job asked for, empty meaning any. Chapter 15A.
133728e barerepo 1mo
27
Labels []string
133728e barerepo 1mo
28
// Name is the workflow job and its matrix combination, so two builds of one commit are told apart.
133728e barerepo 1mo
29
Name string
133728e barerepo 1mo
30
State string
133728e barerepo 1mo
31
RunnerID int64
133728e barerepo 1mo
32
Attempts int
133728e barerepo 1mo
33
Created time.Time
133728e barerepo 1mo
34
Started time.Time
133728e barerepo 1mo
35
Log string
133728e barerepo 1mo
36
}
133728e barerepo 1mo
37
133728e barerepo 1mo
38
// QueueJob adds a build to a repository's queue.
133728e barerepo 1mo
39
func (db *DB) QueueJob(ctx context.Context, repo, ref, sha, command, image string) (int64, error) {
133728e barerepo 1mo
40
return db.QueueJobFor(ctx, repo, ref, sha, command, image, nil, "")
133728e barerepo 1mo
41
}
133728e barerepo 1mo
42
133728e barerepo 1mo
43
// QueueJobFor queues a job only a machine with these labels may take, under a name. Chapter 15A.
133728e barerepo 1mo
44
func (db *DB) QueueJobFor(ctx context.Context, repo, ref, sha, command, image string, labels []string, name string) (int64, error) {
133728e barerepo 1mo
45
var id int64
133728e barerepo 1mo
46
err := db.QueryRowContext(ctx,
133728e barerepo 1mo
47
`INSERT INTO jobs (repo, ref, sha, command, image, labels, name, state, created_at)
133728e barerepo 1mo
48
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?) RETURNING id`,
133728e barerepo 1mo
49
repo, ref, sha, command, image, strings.Join(labels, ","), name,
133728e barerepo 1mo
50
JobQueued, now().Unix()).Scan(&id)
133728e barerepo 1mo
51
return id, err
133728e barerepo 1mo
52
}
133728e barerepo 1mo
53
133728e barerepo 1mo
54
// TakeJob hands over the oldest queued job, first in and first out. Chapter 15.
133728e barerepo 1mo
55
func (db *DB) TakeJob(ctx context.Context, repo string, runner Runner) (*Job, error) {
133728e barerepo 1mo
56
tx, err := db.Begin(ctx)
133728e barerepo 1mo
57
if err != nil {
133728e barerepo 1mo
58
return nil, err
133728e barerepo 1mo
59
}
133728e barerepo 1mo
60
defer tx.Rollback()
133728e barerepo 1mo
61
133728e barerepo 1mo
62
// The first job this machine can run is taken, so one waiting for another machine blocks nothing.
133728e barerepo 1mo
63
rows, err := tx.QueryContext(ctx,
133728e barerepo 1mo
64
`SELECT id, repo, ref, sha, command, image, labels, name, attempts, created_at
133728e barerepo 1mo
65
FROM jobs WHERE repo = ? AND state = ? ORDER BY id`, repo, JobQueued)
133728e barerepo 1mo
66
if err != nil {
133728e barerepo 1mo
67
return nil, err
133728e barerepo 1mo
68
}
133728e barerepo 1mo
69
var j Job
133728e barerepo 1mo
70
var created int64
133728e barerepo 1mo
71
found := false
133728e barerepo 1mo
72
for rows.Next() {
133728e barerepo 1mo
73
var cand Job
133728e barerepo 1mo
74
var labels string
133728e barerepo 1mo
75
var at int64
133728e barerepo 1mo
76
if err := rows.Scan(&cand.ID, &cand.Repo, &cand.Ref, &cand.SHA,
133728e barerepo 1mo
77
&cand.Command, &cand.Image, &labels, &cand.Name, &cand.Attempts, &at); err != nil {
133728e barerepo 1mo
78
rows.Close()
133728e barerepo 1mo
79
return nil, err
133728e barerepo 1mo
80
}
133728e barerepo 1mo
81
cand.Labels = splitLabels(labels)
133728e barerepo 1mo
82
if !runner.Satisfies(cand.Labels) {
133728e barerepo 1mo
83
continue
133728e barerepo 1mo
84
}
133728e barerepo 1mo
85
j, created, found = cand, at, true
133728e barerepo 1mo
86
break
133728e barerepo 1mo
87
}
133728e barerepo 1mo
88
rows.Close()
133728e barerepo 1mo
89
if err := rows.Err(); err != nil {
133728e barerepo 1mo
90
return nil, err
133728e barerepo 1mo
91
}
133728e barerepo 1mo
92
if !found {
133728e barerepo 1mo
93
return nil, nil
133728e barerepo 1mo
94
}
133728e barerepo 1mo
95
j.Created = time.Unix(created, 0)
133728e barerepo 1mo
96
j.State = JobRunning
133728e barerepo 1mo
97
j.RunnerID = runner.ID
133728e barerepo 1mo
98
j.Started = now()
133728e barerepo 1mo
99
133728e barerepo 1mo
100
// Still queued is part of the write, so two machines reaching for one job leave one holding it.
133728e barerepo 1mo
101
res, err := tx.ExecContext(ctx,
133728e barerepo 1mo
102
`UPDATE jobs SET state = ?, runner_id = ?, started_at = ?, attempts = attempts + 1
133728e barerepo 1mo
103
WHERE id = ? AND state = ?`,
133728e barerepo 1mo
104
JobRunning, runner.ID, j.Started.Unix(), j.ID, JobQueued)
133728e barerepo 1mo
105
if err != nil {
133728e barerepo 1mo
106
return nil, err
133728e barerepo 1mo
107
}
133728e barerepo 1mo
108
took, err := res.RowsAffected()
133728e barerepo 1mo
109
if err != nil {
133728e barerepo 1mo
110
return nil, err
133728e barerepo 1mo
111
}
133728e barerepo 1mo
112
if took == 0 {
133728e barerepo 1mo
113
// Somebody else took it between the read and the write, and the runner asks again on its tick.
133728e barerepo 1mo
114
return nil, nil
133728e barerepo 1mo
115
}
133728e barerepo 1mo
116
if err := tx.Commit(); err != nil {
133728e barerepo 1mo
117
return nil, err
133728e barerepo 1mo
118
}
133728e barerepo 1mo
119
return &j, nil
133728e barerepo 1mo
120
}
133728e barerepo 1mo
121
133728e barerepo 1mo
122
// Job reads one job.
133728e barerepo 1mo
123
func (db *DB) Job(ctx context.Context, id int64) (*Job, error) {
133728e barerepo 1mo
124
var j Job
133728e barerepo 1mo
125
var runner sql.NullInt64
133728e barerepo 1mo
126
var created int64
133728e barerepo 1mo
127
var started sql.NullInt64
133728e barerepo 1mo
128
err := db.QueryRowContext(ctx,
133728e barerepo 1mo
129
`SELECT id, repo, ref, sha, command, image, name, state, runner_id, attempts,
133728e barerepo 1mo
130
created_at, started_at, log
133728e barerepo 1mo
131
FROM jobs WHERE id = ?`, id).
133728e barerepo 1mo
132
Scan(&j.ID, &j.Repo, &j.Ref, &j.SHA, &j.Command, &j.Image, &j.Name, &j.State,
133728e barerepo 1mo
133
&runner, &j.Attempts, &created, &started, &j.Log)
133728e barerepo 1mo
134
if errors.Is(err, sql.ErrNoRows) {
133728e barerepo 1mo
135
return nil, ErrNotFound
133728e barerepo 1mo
136
}
133728e barerepo 1mo
137
if err != nil {
133728e barerepo 1mo
138
return nil, err
133728e barerepo 1mo
139
}
133728e barerepo 1mo
140
j.RunnerID = runner.Int64
133728e barerepo 1mo
141
j.Created = time.Unix(created, 0)
133728e barerepo 1mo
142
if started.Valid {
133728e barerepo 1mo
143
j.Started = time.Unix(started.Int64, 0)
133728e barerepo 1mo
144
}
133728e barerepo 1mo
145
return &j, nil
133728e barerepo 1mo
146
}
133728e barerepo 1mo
147
133728e barerepo 1mo
148
// AppendLog adds a chunk of build output to a running job.
133728e barerepo 1mo
149
func (db *DB) AppendLog(ctx context.Context, id int64, runnerID int64, chunk string) error {
133728e barerepo 1mo
150
res, err := db.ExecContext(ctx,
133728e barerepo 1mo
151
`UPDATE jobs SET log = log || ? WHERE id = ? AND runner_id = ? AND state = ?`,
133728e barerepo 1mo
152
chunk, id, runnerID, JobRunning)
133728e barerepo 1mo
153
if err != nil {
133728e barerepo 1mo
154
return err
133728e barerepo 1mo
155
}
133728e barerepo 1mo
156
if n, _ := res.RowsAffected(); n == 0 {
133728e barerepo 1mo
157
return ErrNotFound
133728e barerepo 1mo
158
}
133728e barerepo 1mo
159
return nil
133728e barerepo 1mo
160
}
133728e barerepo 1mo
161
133728e barerepo 1mo
162
// FinishJob closes a job out and returns it, so the caller can write the note.
133728e barerepo 1mo
163
func (db *DB) FinishJob(ctx context.Context, id, runnerID int64) (*Job, error) {
133728e barerepo 1mo
164
res, err := db.ExecContext(ctx,
133728e barerepo 1mo
165
`UPDATE jobs SET state = ? WHERE id = ? AND runner_id = ? AND state = ?`,
133728e barerepo 1mo
166
JobDone, id, runnerID, JobRunning)
133728e barerepo 1mo
167
if err != nil {
133728e barerepo 1mo
168
return nil, err
133728e barerepo 1mo
169
}
133728e barerepo 1mo
170
if n, _ := res.RowsAffected(); n == 0 {
133728e barerepo 1mo
171
return nil, ErrNotFound
133728e barerepo 1mo
172
}
133728e barerepo 1mo
173
return db.Job(ctx, id)
133728e barerepo 1mo
174
}
133728e barerepo 1mo
175
133728e barerepo 1mo
176
// RequeueLostJobs retries a silent runner's job once, then marks it lost rather than cycling.
133728e barerepo 1mo
177
func (db *DB) RequeueLostJobs(ctx context.Context, olderThan time.Duration) (int, error) {
133728e barerepo 1mo
178
cutoff := now().Add(-olderThan).Unix()
133728e barerepo 1mo
179
res, err := db.ExecContext(ctx,
133728e barerepo 1mo
180
`UPDATE jobs SET state = ?, runner_id = NULL
133728e barerepo 1mo
181
WHERE state = ? AND started_at < ? AND attempts < 2`,
133728e barerepo 1mo
182
JobQueued, JobRunning, cutoff)
133728e barerepo 1mo
183
if err != nil {
133728e barerepo 1mo
184
return 0, err
133728e barerepo 1mo
185
}
133728e barerepo 1mo
186
requeued, _ := res.RowsAffected()
133728e barerepo 1mo
187
133728e barerepo 1mo
188
if _, err := db.ExecContext(ctx,
133728e barerepo 1mo
189
`UPDATE jobs SET state = ? WHERE state = ? AND started_at < ? AND attempts >= 2`,
133728e barerepo 1mo
190
JobLost, JobRunning, cutoff); err != nil {
133728e barerepo 1mo
191
return int(requeued), err
133728e barerepo 1mo
192
}
133728e barerepo 1mo
193
return int(requeued), nil
133728e barerepo 1mo
194
}
133728e barerepo 1mo
195
133728e barerepo 1mo
196
// QueuedJobs lists what is waiting for a machine, which the runs page and the tests both ask for.
133728e barerepo 1mo
197
func (db *DB) QueuedJobs(ctx context.Context, repo string) ([]Job, error) {
133728e barerepo 1mo
198
rows, err := db.QueryContext(ctx,
133728e barerepo 1mo
199
`SELECT id, repo, ref, sha, command, image, labels, name, attempts, created_at
133728e barerepo 1mo
200
FROM jobs WHERE repo = ? AND state = ? ORDER BY id`, repo, JobQueued)
133728e barerepo 1mo
201
if err != nil {
133728e barerepo 1mo
202
return nil, err
133728e barerepo 1mo
203
}
133728e barerepo 1mo
204
defer rows.Close()
133728e barerepo 1mo
205
var out []Job
133728e barerepo 1mo
206
for rows.Next() {
133728e barerepo 1mo
207
var j Job
133728e barerepo 1mo
208
var labels string
133728e barerepo 1mo
209
var created int64
133728e barerepo 1mo
210
if err := rows.Scan(&j.ID, &j.Repo, &j.Ref, &j.SHA,
133728e barerepo 1mo
211
&j.Command, &j.Image, &labels, &j.Name, &j.Attempts, &created); err != nil {
133728e barerepo 1mo
212
return nil, err
133728e barerepo 1mo
213
}
133728e barerepo 1mo
214
j.Labels = splitLabels(labels)
133728e barerepo 1mo
215
j.State = JobQueued
133728e barerepo 1mo
216
j.Created = time.Unix(created, 0)
133728e barerepo 1mo
217
out = append(out, j)
133728e barerepo 1mo
218
}
133728e barerepo 1mo
219
return out, rows.Err()
133728e barerepo 1mo
220
}
history · rawbarerepo 0.1.0