internal/store/queues.go

a32f7f2001c1c0bf38ba8c3360e6bc7ca3982fee
gitbay/internal/store/queues.go history · blame · raw

226 lines · 8422 bytes

  1package store
  2
  3import "fmt"
  4
  5// Queues is the state of every background worker, for the instance admin:
  6// what is waiting, what is retrying, and what has given up. Item lists are
  7// capped so a flood of one kind cannot bury the others.
  8type Queues struct {
  9	Webhooks QueueWebhooks `json:"webhooks"`
 10	Mail     QueueMail     `json:"mail"`
 11	Mirrors  QueueMirrors  `json:"mirrors"`
 12	Builds   QueueBuilds   `json:"builds"`
 13	Deps     QueueDeps     `json:"deps"`
 14}
 15
 16type QueueWebhooks struct {
 17	Pending       int64              `json:"pending"`
 18	Retrying      int64              `json:"retrying"` // pending with at least one failed attempt
 19	Failed        int64              `json:"failed"`   // dead-lettered
 20	OldestPending string             `json:"oldest_pending,omitempty"`
 21	Items         []QueueDeliveryRow `json:"items"` // retrying and dead-lettered, newest first
 22}
 23
 24type QueueDeliveryRow struct {
 25	ID        int64  `json:"id"`
 26	Repo      string `json:"repo"`
 27	URL       string `json:"url"`
 28	Attempts  int64  `json:"attempts"`
 29	Status    int64  `json:"last_status,omitempty"`
 30	LastError string `json:"last_error,omitempty"`
 31	FailedAt  string `json:"failed_at,omitempty"`
 32	CreatedAt string `json:"created_at"`
 33}
 34
 35type QueueMail struct {
 36	Pending       int64          `json:"pending"`
 37	Retrying      int64          `json:"retrying"`
 38	Failed        int64          `json:"failed"`
 39	OldestPending string         `json:"oldest_pending,omitempty"`
 40	Items         []QueueMailRow `json:"items"`
 41}
 42
 43type QueueMailRow struct {
 44	ID        int64  `json:"id"`
 45	Recipient string `json:"recipient"`
 46	Subject   string `json:"subject"`
 47	Attempts  int64  `json:"attempts"`
 48	LastError string `json:"last_error,omitempty"`
 49	FailedAt  string `json:"failed_at,omitempty"`
 50	CreatedAt string `json:"created_at"`
 51}
 52
 53type QueueMirrors struct {
 54	Dirty  int64            `json:"dirty"` // waiting for a sync
 55	Errors int64            `json:"errors"`
 56	Items  []QueueMirrorRow `json:"items"` // the ones whose last sync failed
 57}
 58
 59type QueueMirrorRow struct {
 60	ID        int64  `json:"id"`
 61	Repo      string `json:"repo"`
 62	Direction string `json:"direction"`
 63	URL       string `json:"url"`
 64	LastSync  string `json:"last_sync,omitempty"`
 65	LastError string `json:"last_error"`
 66}
 67
 68type QueueBuilds struct {
 69	Pending       int64           `json:"pending"`
 70	Running       int64           `json:"running"`
 71	OldestPending string          `json:"oldest_pending,omitempty"`
 72	Items         []QueueBuildRow `json:"items"` // running then pending, oldest first
 73}
 74
 75type QueueBuildRow struct {
 76	Repo      string `json:"repo"`
 77	Number    int64  `json:"number"`
 78	Job       string `json:"job"`
 79	Status    string `json:"status"` // running | pending
 80	CreatedAt string `json:"created_at"`
 81	StartedAt string `json:"started_at"` // "" while pending
 82}
 83
 84type QueueDeps struct {
 85	Errors int64         `json:"errors"`
 86	Items  []QueueDepRow `json:"items"`
 87}
 88
 89type QueueDepRow struct {
 90	Repo      string `json:"repo"`
 91	LastCheck string `json:"last_check,omitempty"`
 92	LastError string `json:"last_error"`
 93}
 94
 95const queueItemCap = 20
 96
 97const repoPathExpr = `COALESCE(u.username, o.name) || '/' || r.name`
 98
 99// repoJoin joins repos and their owner for the path expression; the
100// argument is the column holding the repo id.
101func repoJoin(col string) string {
102	return fmt.Sprintf(` JOIN repos r ON r.id = %s
103		LEFT JOIN users u ON r.owner_kind = 'user' AND u.id = r.owner_id
104		LEFT JOIN orgs o ON r.owner_kind = 'org' AND o.id = r.owner_id`, col)
105}
106
107// QueueStatus reads every worker queue. Read-only; safe on a live daemon.
108func (s *Store) QueueStatus() (Queues, error) {
109	q := Queues{
110		Webhooks: QueueWebhooks{Items: []QueueDeliveryRow{}},
111		Mail:     QueueMail{Items: []QueueMailRow{}},
112		Mirrors:  QueueMirrors{Items: []QueueMirrorRow{}},
113		Builds:   QueueBuilds{Items: []QueueBuildRow{}},
114		Deps:     QueueDeps{Items: []QueueDepRow{}},
115	}
116
117	if err := s.DB.QueryRow(`SELECT
118		COUNT(*) FILTER (WHERE delivered_at IS NULL AND failed_at IS NULL),
119		COUNT(*) FILTER (WHERE delivered_at IS NULL AND failed_at IS NULL AND attempts > 0),
120		COUNT(*) FILTER (WHERE failed_at IS NOT NULL),
121		COALESCE(MIN(created_at) FILTER (WHERE delivered_at IS NULL AND failed_at IS NULL), '')
122		FROM webhook_deliveries`).Scan(&q.Webhooks.Pending, &q.Webhooks.Retrying, &q.Webhooks.Failed, &q.Webhooks.OldestPending); err != nil {
123		return q, err
124	}
125	if err := s.queryEach(`SELECT d.id, `+repoPathExpr+`, w.url, d.attempts, COALESCE(d.last_status, 0),
126		COALESCE(d.last_error, ''), COALESCE(d.failed_at, ''), d.created_at
127		FROM webhook_deliveries d JOIN webhooks w ON w.id = d.webhook_id`+repoJoin("w.repo_id")+`
128		WHERE d.delivered_at IS NULL AND (d.failed_at IS NOT NULL OR d.attempts > 0)
129		ORDER BY d.id DESC LIMIT ?`, func(sc scanner) error {
130		var d QueueDeliveryRow
131		if err := sc.Scan(&d.ID, &d.Repo, &d.URL, &d.Attempts, &d.Status, &d.LastError, &d.FailedAt, &d.CreatedAt); err != nil {
132			return err
133		}
134		q.Webhooks.Items = append(q.Webhooks.Items, d)
135		return nil
136	}); err != nil {
137		return q, err
138	}
139
140	if err := s.DB.QueryRow(`SELECT
141		COUNT(*) FILTER (WHERE sent_at IS NULL AND failed_at IS NULL),
142		COUNT(*) FILTER (WHERE sent_at IS NULL AND failed_at IS NULL AND attempts > 0),
143		COUNT(*) FILTER (WHERE failed_at IS NOT NULL),
144		COALESCE(MIN(created_at) FILTER (WHERE sent_at IS NULL AND failed_at IS NULL), '')
145		FROM notifications`).Scan(&q.Mail.Pending, &q.Mail.Retrying, &q.Mail.Failed, &q.Mail.OldestPending); err != nil {
146		return q, err
147	}
148	if err := s.queryEach(`SELECT id, recipient, subject, attempts, COALESCE(last_error, ''), COALESCE(failed_at, ''), created_at
149		FROM notifications WHERE sent_at IS NULL AND (failed_at IS NOT NULL OR attempts > 0)
150		ORDER BY id DESC LIMIT ?`, func(sc scanner) error {
151		var m QueueMailRow
152		if err := sc.Scan(&m.ID, &m.Recipient, &m.Subject, &m.Attempts, &m.LastError, &m.FailedAt, &m.CreatedAt); err != nil {
153			return err
154		}
155		q.Mail.Items = append(q.Mail.Items, m)
156		return nil
157	}); err != nil {
158		return q, err
159	}
160
161	if err := s.DB.QueryRow(`SELECT COUNT(*) FILTER (WHERE dirty = 1), COUNT(*) FILTER (WHERE last_error != '')
162		FROM mirrors`).Scan(&q.Mirrors.Dirty, &q.Mirrors.Errors); err != nil {
163		return q, err
164	}
165	if err := s.queryEach(`SELECT m.id, `+repoPathExpr+`, m.direction, m.url, m.last_sync, m.last_error
166		FROM mirrors m`+repoJoin("m.repo_id")+` WHERE m.last_error != '' ORDER BY m.id DESC LIMIT ?`, func(sc scanner) error {
167		var m QueueMirrorRow
168		if err := sc.Scan(&m.ID, &m.Repo, &m.Direction, &m.URL, &m.LastSync, &m.LastError); err != nil {
169			return err
170		}
171		q.Mirrors.Items = append(q.Mirrors.Items, m)
172		return nil
173	}); err != nil {
174		return q, err
175	}
176
177	if err := s.DB.QueryRow(`SELECT COUNT(*) FILTER (WHERE status = 'pending'), COUNT(*) FILTER (WHERE status = 'running'),
178		COALESCE(MIN(created_at) FILTER (WHERE status = 'pending'), '') FROM builds`).Scan(&q.Builds.Pending, &q.Builds.Running, &q.Builds.OldestPending); err != nil {
179		return q, err
180	}
181	if err := s.queryEach(`SELECT `+repoPathExpr+`, b.number, b.job, b.status, b.created_at, b.started_at
182		FROM builds b`+repoJoin("b.repo_id")+` WHERE b.status IN ('running', 'pending')
183		ORDER BY b.status = 'running' DESC, b.started_at, b.created_at, b.id LIMIT ?`, func(sc scanner) error {
184		var b QueueBuildRow
185		if err := sc.Scan(&b.Repo, &b.Number, &b.Job, &b.Status, &b.CreatedAt, &b.StartedAt); err != nil {
186			return err
187		}
188		q.Builds.Items = append(q.Builds.Items, b)
189		return nil
190	}); err != nil {
191		return q, err
192	}
193
194	if err := s.DB.QueryRow(`SELECT COUNT(*) FROM dep_checks WHERE last_error != ''`).Scan(&q.Deps.Errors); err != nil {
195		return q, err
196	}
197	if err := s.queryEach(`SELECT `+repoPathExpr+`, c.last_check, c.last_error
198		FROM dep_checks c`+repoJoin("c.repo_id")+` WHERE c.last_error != '' ORDER BY c.repo_id LIMIT ?`, func(sc scanner) error {
199		var d QueueDepRow
200		if err := sc.Scan(&d.Repo, &d.LastCheck, &d.LastError); err != nil {
201			return err
202		}
203		q.Deps.Items = append(q.Deps.Items, d)
204		return nil
205	}); err != nil {
206		return q, err
207	}
208	return q, nil
209}
210
211type scanner interface{ Scan(...any) error }
212
213// queryEach runs a capped item query and hands each row to fn.
214func (s *Store) queryEach(query string, fn func(scanner) error) error {
215	rows, err := s.DB.Query(query, queueItemCap)
216	if err != nil {
217		return err
218	}
219	defer rows.Close()
220	for rows.Next() {
221		if err := fn(rows); err != nil {
222			return err
223		}
224	}
225	return rows.Err()
226}