internal/store/queues.go

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

269 lines · 10184 bytes

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