internal/store/builds.go

477 lines · 15624 bytes

  1package store
  2
  3import (
  4	"database/sql"
  5	"errors"
  6	"os"
  7	"strconv"
  8	"strings"
  9	"time"
 10)
 11
 12// Build is one CI job execution for one commit.
 13type Build struct {
 14	ID         int64
 15	RepoID     int64
 16	Number     int64
 17	Job        string
 18	SHA        string
 19	Ref        string
 20	Steps      string // JSON array of shell commands
 21	Image      string // container image for the steps; "" means the runner default
 22	Tree       string // the commit's tree; "" when not deduplicated by tree
 23	Status     string // pending|running|success|failure
 24	CreatedAt  string
 25	StartedAt  string
 26	FinishedAt string
 27	// LogClosedAt is when the runner's log stream ended; "" while it is
 28	// open or was never opened. Set on a running build only.
 29	LogClosedAt string
 30	// Trusted is false for a merge request head fetched from another
 31	// repository: its steps run without the target's secrets.
 32	Trusted bool
 33}
 34
 35// MaxBuildLog caps a build's stored log; appends past it are dropped.
 36const MaxBuildLog = 2 << 20
 37
 38// truncNotice is appended once when a log first hits the cap. A log that
 39// simply stops is indistinguishable from a build that died mid-step, which
 40// is the reading that sent people hunting for a nonexistent test failure.
 41var truncNotice = []byte("\n[log truncated: reached the " +
 42	strconv.Itoa(MaxBuildLog>>20) + " MiB cap; earlier output is above]\n")
 43
 44// CreateBuild allocates the per-repo build number in the same transaction
 45// as the insert, like issue and MR numbers.
 46func (s *Store) CreateBuild(repoID int64, job, sha, ref, stepsJSON, image, tree string, trusted bool) (int64, error) {
 47	tx, err := s.DB.Begin()
 48	if err != nil {
 49		return 0, err
 50	}
 51	defer tx.Rollback()
 52	if _, err := tx.Exec("UPDATE repos SET build_counter = build_counter + 1 WHERE id = ?", repoID); err != nil {
 53		return 0, err
 54	}
 55	var n int64
 56	if err := tx.QueryRow("SELECT build_counter FROM repos WHERE id = ?", repoID).Scan(&n); err != nil {
 57		return 0, err
 58	}
 59	if _, err := tx.Exec(
 60		"INSERT INTO builds (repo_id, number, job, sha, ref, steps, image, tree, trusted) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)",
 61		repoID, n, job, sha, ref, stepsJSON, image, tree, trusted); err != nil {
 62		return 0, err
 63	}
 64	return n, tx.Commit()
 65}
 66
 67const buildSelect = `
 68	SELECT id, repo_id, number, job, sha, ref, steps, image, tree, status, created_at, started_at, finished_at, log_closed_at, trusted
 69	FROM builds`
 70
 71func scanBuild(row interface{ Scan(...any) error }) (Build, error) {
 72	var b Build
 73	var trusted int
 74	err := row.Scan(&b.ID, &b.RepoID, &b.Number, &b.Job, &b.SHA, &b.Ref, &b.Steps, &b.Image, &b.Tree,
 75		&b.Status, &b.CreatedAt, &b.StartedAt, &b.FinishedAt, &b.LogClosedAt, &trusted)
 76	b.Trusted = trusted != 0
 77	return b, err
 78}
 79
 80// ClaimBuild atomically hands the oldest pending build to a runner and
 81// marks it running. A non-empty repoIDs restricts the claim to those
 82// repositories. Untrusted builds — merge request heads from another
 83// repository — are skipped unless untrusted is set: they run a stranger's
 84// code, which only a runner that isolates should take.
 85func (s *Store) ClaimBuild(repoIDs []int64, untrusted bool) (Build, bool, error) {
 86	tx, err := s.DB.Begin()
 87	if err != nil {
 88		return Build{}, false, err
 89	}
 90	defer tx.Rollback()
 91	query := "SELECT id FROM builds WHERE status = 'pending'"
 92	args := []any{}
 93	if !untrusted {
 94		query += " AND trusted = 1"
 95	}
 96	if len(repoIDs) > 0 {
 97		marks := strings.TrimSuffix(strings.Repeat("?,", len(repoIDs)), ",")
 98		query += " AND repo_id IN (" + marks + ")"
 99		for _, id := range repoIDs {
100			args = append(args, id)
101		}
102	}
103	query += " ORDER BY id LIMIT 1"
104	var id int64
105	err = tx.QueryRow(query, args...).Scan(&id)
106	if errors.Is(err, sql.ErrNoRows) {
107		return Build{}, false, nil
108	}
109	if err != nil {
110		return Build{}, false, err
111	}
112	if _, err := tx.Exec(
113		"UPDATE builds SET status = 'running', started_at = strftime('%Y-%m-%dT%H:%M:%SZ','now') WHERE id = ?", id); err != nil {
114		return Build{}, false, err
115	}
116	b, err := scanBuild(tx.QueryRow(buildSelect+" WHERE id = ?", id))
117	if err != nil {
118		return Build{}, false, err
119	}
120	return b, true, tx.Commit()
121}
122
123// StaleBuildDeadline is how long a claimed build may stay running before the
124// server gives up on it. Comfortably longer than the runner's own -timeout
125// (45m by default), so this only fires when the runner never reported at all —
126// it was killed, restarted, or lost the network mid-build.
127const StaleBuildDeadline = 90 * time.Minute
128
129// staleBuildDeadline is StaleBuildDeadline unless GITBAY_STALE_BUILD_DEADLINE
130// shortens it, which tests do.
131func staleBuildDeadline() time.Duration {
132	if v := os.Getenv("GITBAY_STALE_BUILD_DEADLINE"); v != "" {
133		if d, err := time.ParseDuration(v); err == nil && d > 0 {
134			return d
135		}
136	}
137	return StaleBuildDeadline
138}
139
140// StaleLogGrace is how long a running build may go on after its log
141// stream ended before it is treated as abandoned. The runner reports the
142// outcome right after closing the stream, retrying for about thirty
143// seconds if the server is unreachable; two minutes outlasts that.
144const StaleLogGrace = 2 * time.Minute
145
146// MarkBuildLogClosed records that the runner's log stream for a build
147// ended, on a build still running. A build that finishes normally is
148// reported moments later and the mark is moot; one that is not has lost
149// its runner, and ReapStaleBuilds fails it after StaleLogGrace rather
150// than at the deadline (#179).
151func (s *Store) MarkBuildLogClosed(id int64) error {
152	_, err := s.DB.Exec(`
153		UPDATE builds SET log_closed_at = strftime('%Y-%m-%dT%H:%M:%SZ','now')
154		WHERE id = ? AND status = 'running' AND log_closed_at = ''`, id)
155	return err
156}
157
158// ReapStaleBuilds fails every running build whose runner is gone and
159// returns them, so the caller can resolve their commit statuses: one whose
160// log stream ended more than StaleLogGrace ago with no outcome reported,
161// or one running past the deadline with no stream ever seen. A runner that
162// dies between claiming a build and reporting it otherwise leaves the row
163// claimed forever, and the commit pending forever with it.
164func (s *Store) ReapStaleBuilds() ([]Build, error) {
165	const layout = "2006-01-02T15:04:05Z"
166	now := time.Now().UTC()
167	cutoff := now.Add(-staleBuildDeadline()).Format(layout)
168	logCutoff := now.Add(-StaleLogGrace).Format(layout)
169	rows, err := s.DB.Query(buildSelect+
170		" WHERE status = 'running' AND ((started_at != '' AND started_at < ?)"+
171		" OR (log_closed_at != '' AND log_closed_at < ?))", cutoff, logCutoff)
172	if err != nil {
173		return nil, err
174	}
175	defer rows.Close()
176	var stale []Build
177	for rows.Next() {
178		b, err := scanBuild(rows)
179		if err != nil {
180			return nil, err
181		}
182		stale = append(stale, b)
183	}
184	if err := rows.Err(); err != nil {
185		return nil, err
186	}
187	for _, b := range stale {
188		if err := s.AppendBuildLog(b.ID, []byte(
189			"\nbuild abandoned: the runner never reported an outcome\n")); err != nil {
190			return nil, err
191		}
192		if err := s.FinishBuild(b.ID, "failure"); err != nil {
193			return nil, err
194		}
195		if _, err := s.DB.Exec(`UPDATE builds SET reaped_at = finished_at WHERE id = ?`, b.ID); err != nil {
196			return nil, err
197		}
198	}
199	return stale, nil
200}
201
202// QueueStats is the state of the build queue: what waits now, and over
203// the last day how long a build waited to be claimed and how many were
204// ended by the reaper rather than by a runner's report (#184).
205type QueueStats struct {
206	Pending       int64 `json:"pending"`
207	Claimed24h    int64 `json:"claimed_24h"`
208	ClaimWaitAvgS int64 `json:"claim_wait_avg_s"`
209	ClaimWaitMaxS int64 `json:"claim_wait_max_s"`
210	Reaped24h     int64 `json:"reaped_24h"`
211}
212
213func (s *Store) QueueStats() (QueueStats, error) {
214	var q QueueStats
215	since := time.Now().UTC().Add(-24 * time.Hour).Format("2006-01-02T15:04:05Z")
216	err := s.DB.QueryRow(`SELECT
217		(SELECT COUNT(*) FROM builds WHERE status = 'pending'),
218		COUNT(*),
219		CAST(COALESCE(AVG(strftime('%s', started_at) - strftime('%s', created_at)), 0) AS INTEGER),
220		CAST(COALESCE(MAX(strftime('%s', started_at) - strftime('%s', created_at)), 0) AS INTEGER),
221		(SELECT COUNT(*) FROM builds WHERE reaped_at >= ?)
222		FROM builds WHERE started_at >= ?`, since, since).
223		Scan(&q.Pending, &q.Claimed24h, &q.ClaimWaitAvgS, &q.ClaimWaitMaxS, &q.Reaped24h)
224	return q, err
225}
226
227// AppendBuildLog adds a chunk to the build's log, dropping bytes past the cap.
228func (s *Store) AppendBuildLog(id int64, chunk []byte) error {
229	res, err := s.DB.Exec(`
230		UPDATE builds SET log = log || ?
231		WHERE id = ? AND length(log) < ?`, chunk, id, MaxBuildLog)
232	if err != nil {
233		return err
234	}
235	if n, _ := res.RowsAffected(); n > 0 {
236		s.wakeBuild(id)
237		return nil
238	}
239	// Over the cap. The bounds match exactly once: appending the notice puts
240	// the log past the upper bound, so later chunks fall through silently.
241	_, err = s.DB.Exec(`
242		UPDATE builds SET log = log || ?
243		WHERE id = ? AND length(log) >= ? AND length(log) < ?`,
244		truncNotice, id, MaxBuildLog, MaxBuildLog+len(truncNotice))
245	if err == nil {
246		s.wakeBuild(id)
247	}
248	return err
249}
250
251// FinishBuild records the outcome of a running build.
252func (s *Store) FinishBuild(id int64, status string) error {
253	res, err := s.DB.Exec(`
254		UPDATE builds SET status = ?, finished_at = strftime('%Y-%m-%dT%H:%M:%SZ','now')
255		WHERE id = ? AND status = 'running'`, status, id)
256	if err != nil {
257		return err
258	}
259	if n, _ := res.RowsAffected(); n == 0 {
260		return ErrNotFound
261	}
262	s.wakeBuild(id)
263	return nil
264}
265
266func (s *Store) BuildByID(id int64) (Build, error) {
267	b, err := scanBuild(s.DB.QueryRow(buildSelect+" WHERE id = ?", id))
268	if errors.Is(err, sql.ErrNoRows) {
269		return b, ErrNotFound
270	}
271	return b, err
272}
273
274func (s *Store) BuildByNumber(repoID, number int64) (Build, error) {
275	b, err := scanBuild(s.DB.QueryRow(buildSelect+" WHERE repo_id = ? AND number = ?", repoID, number))
276	if errors.Is(err, sql.ErrNoRows) {
277		return b, ErrNotFound
278	}
279	return b, err
280}
281
282// BuildFilter narrows ListBuilds to builds matching every non-empty field.
283// Before is the keyset cursor: only builds numbered below it, which with
284// the newest-first order is the page after the one that ended there.
285type BuildFilter struct {
286	Ref    string
287	Status string
288	Job    string
289	Before int64
290}
291
292func (s *Store) ListBuilds(repoID int64, f BuildFilter, limit int) ([]Build, error) {
293	q := buildSelect + " WHERE repo_id = ?"
294	args := []any{repoID}
295	if f.Ref != "" {
296		q += " AND ref = ?"
297		args = append(args, f.Ref)
298	}
299	if f.Status != "" {
300		q += " AND status = ?"
301		args = append(args, f.Status)
302	}
303	if f.Job != "" {
304		q += " AND job = ?"
305		args = append(args, f.Job)
306	}
307	if f.Before > 0 {
308		q += " AND number < ?"
309		args = append(args, f.Before)
310	}
311	q += " ORDER BY number DESC LIMIT ?"
312	args = append(args, limit)
313	rows, err := s.DB.Query(q, args...)
314	if err != nil {
315		return nil, err
316	}
317	defer rows.Close()
318	var out []Build
319	for rows.Next() {
320		b, err := scanBuild(rows)
321		if err != nil {
322			return nil, err
323		}
324		out = append(out, b)
325	}
326	return out, rows.Err()
327}
328
329// BuildLog returns the stored log bytes.
330func (s *Store) BuildLog(id int64) ([]byte, error) {
331	var log []byte
332	err := s.DB.QueryRow("SELECT log FROM builds WHERE id = ?", id).Scan(&log)
333	if errors.Is(err, sql.ErrNoRows) {
334		return nil, ErrNotFound
335	}
336	return log, err
337}
338
339// BuildLogWait returns a channel closed by the next append to, finish of
340// or cancel of the build. Take it before reading, so a change between the
341// read and the wait still wakes the reader. Only this process's writes
342// wake it.
343func (s *Store) BuildLogWait(id int64) <-chan struct{} {
344	s.logMu.Lock()
345	defer s.logMu.Unlock()
346	if s.logWait == nil {
347		s.logWait = map[int64]chan struct{}{}
348	}
349	ch, ok := s.logWait[id]
350	if !ok {
351		ch = make(chan struct{})
352		s.logWait[id] = ch
353	}
354	return ch
355}
356
357func (s *Store) wakeBuild(id int64) {
358	s.logMu.Lock()
359	defer s.logMu.Unlock()
360	if ch, ok := s.logWait[id]; ok {
361		close(ch)
362		delete(s.logWait, id)
363	}
364}
365
366// BuildLogFrom returns the build's status and its log past offset bytes,
367// read together so a terminal status comes with every byte before it.
368// The cast matters: || stores the log as text, and substr on text counts
369// characters.
370func (s *Store) BuildLogFrom(id, offset int64) (string, []byte, error) {
371	var status string
372	var chunk []byte
373	err := s.DB.QueryRow(`SELECT status, substr(CAST(log AS BLOB), ?) FROM builds WHERE id = ?`,
374		offset+1, id).Scan(&status, &chunk)
375	if errors.Is(err, sql.ErrNoRows) {
376		return "", nil, ErrNotFound
377	}
378	return status, chunk, err
379}
380
381// LatestBuild returns the newest build for a repo, optionally narrowed to
382// one job. It is what a status badge reports.
383func (s *Store) LatestBuild(repoID int64, job string) (Build, error) {
384	q := buildSelect + " WHERE repo_id = ?"
385	args := []any{repoID}
386	if job != "" {
387		q += " AND job = ?"
388		args = append(args, job)
389	}
390	q += " ORDER BY number DESC LIMIT 1"
391	b, err := scanBuild(s.DB.QueryRow(q, args...))
392	if errors.Is(err, sql.ErrNoRows) {
393		return b, ErrNotFound
394	}
395	return b, err
396}
397
398// BuildsForCommit returns the newest build per job for one commit. A merge
399// request's checks are ci/<job> statuses; this is where their timing comes
400// from, in one query rather than one per check.
401func (s *Store) BuildsForCommit(repoID int64, sha string) (map[string]Build, error) {
402	rows, err := s.DB.Query(buildSelect+" WHERE repo_id = ? AND sha = ? ORDER BY number ASC", repoID, sha)
403	if err != nil {
404		return nil, err
405	}
406	defer rows.Close()
407	out := map[string]Build{}
408	for rows.Next() {
409		b, err := scanBuild(rows)
410		if err != nil {
411			return nil, err
412		}
413		out[b.Job] = b // ascending: the last row for a job wins
414	}
415	return out, rows.Err()
416}
417
418// Elapsed reports how long a build ran. Zero until it has both a start and
419// a finish, which is every state but success and failure.
420func (b Build) Elapsed() time.Duration {
421	const layout = "2006-01-02T15:04:05Z"
422	start, err := time.Parse(layout, b.StartedAt)
423	if err != nil {
424		return 0
425	}
426	end, err := time.Parse(layout, b.FinishedAt)
427	if err != nil {
428		return 0
429	}
430	if d := end.Sub(start); d > 0 {
431		return d.Round(time.Second)
432	}
433	return 0
434}
435
436// CancelBuild withdraws a queued or running build. A running one is
437// ended by the runner, which learns of the cancellation when its log
438// session is closed, and whose later report lands on a row that already
439// says cancelled.
440func (s *Store) CancelBuild(id int64) error {
441	res, err := s.DB.Exec(`UPDATE builds SET status = 'cancelled',
442		finished_at = strftime('%Y-%m-%dT%H:%M:%SZ','now') WHERE id = ? AND status IN ('pending', 'running')`, id)
443	if err != nil {
444		return err
445	}
446	if n, _ := res.RowsAffected(); n == 0 {
447		return ErrNotFound
448	}
449	s.wakeBuild(id)
450	return nil
451}
452
453// SuccessBuildFor finds a passed build of the commit for the job, on any
454// ref: what a cancelled duplicate can point back at.
455// SuccessBuildForTree is SuccessBuildFor keyed by tree rather than
456// commit: a rebase that changes nothing in the tree has already been
457// built (#177). An empty tree never matches.
458func (s *Store) SuccessBuildForTree(repoID int64, tree, job string) (Build, bool, error) {
459	if tree == "" {
460		return Build{}, false, nil
461	}
462	b, err := scanBuild(s.DB.QueryRow(buildSelect+
463		" WHERE repo_id = ? AND tree = ? AND job = ? AND status = 'success' ORDER BY number DESC LIMIT 1", repoID, tree, job))
464	if errors.Is(err, sql.ErrNoRows) {
465		return Build{}, false, nil
466	}
467	return b, err == nil, err
468}
469
470func (s *Store) SuccessBuildFor(repoID int64, sha, job string) (Build, bool, error) {
471	b, err := scanBuild(s.DB.QueryRow(buildSelect+
472		" WHERE repo_id = ? AND sha = ? AND job = ? AND status = 'success' ORDER BY number DESC LIMIT 1", repoID, sha, job))
473	if errors.Is(err, sql.ErrNoRows) {
474		return Build{}, false, nil
475	}
476	return b, err == nil, err
477}