internal/store/queues.go

7effe3feb4612cffb5ed38c02300299559aa8971
gitbay/internal/store/queues.go history · blame · raw

223 lines · 8191 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 builds, oldest first
 73}
 74
 75type QueueBuildRow struct {
 76	Repo      string `json:"repo"`
 77	Number    int64  `json:"number"`
 78	Job       string `json:"job"`
 79	StartedAt string `json:"started_at"`
 80}
 81
 82type QueueDeps struct {
 83	Errors int64         `json:"errors"`
 84	Items  []QueueDepRow `json:"items"`
 85}
 86
 87type QueueDepRow struct {
 88	Repo      string `json:"repo"`
 89	LastCheck string `json:"last_check,omitempty"`
 90	LastError string `json:"last_error"`
 91}
 92
 93const queueItemCap = 20
 94
 95const repoPathExpr = `COALESCE(u.username, o.name) || '/' || r.name`
 96
 97// repoJoin joins repos and their owner for the path expression; the
 98// argument is the column holding the repo id.
 99func repoJoin(col string) string {
100	return fmt.Sprintf(` JOIN repos r ON r.id = %s
101		LEFT JOIN users u ON r.owner_kind = 'user' AND u.id = r.owner_id
102		LEFT JOIN orgs o ON r.owner_kind = 'org' AND o.id = r.owner_id`, col)
103}
104
105// QueueStatus reads every worker queue. Read-only; safe on a live daemon.
106func (s *Store) QueueStatus() (Queues, error) {
107	q := Queues{
108		Webhooks: QueueWebhooks{Items: []QueueDeliveryRow{}},
109		Mail:     QueueMail{Items: []QueueMailRow{}},
110		Mirrors:  QueueMirrors{Items: []QueueMirrorRow{}},
111		Builds:   QueueBuilds{Items: []QueueBuildRow{}},
112		Deps:     QueueDeps{Items: []QueueDepRow{}},
113	}
114
115	if err := s.DB.QueryRow(`SELECT
116		COUNT(*) FILTER (WHERE delivered_at IS NULL AND failed_at IS NULL),
117		COUNT(*) FILTER (WHERE delivered_at IS NULL AND failed_at IS NULL AND attempts > 0),
118		COUNT(*) FILTER (WHERE failed_at IS NOT NULL),
119		COALESCE(MIN(created_at) FILTER (WHERE delivered_at IS NULL AND failed_at IS NULL), '')
120		FROM webhook_deliveries`).Scan(&q.Webhooks.Pending, &q.Webhooks.Retrying, &q.Webhooks.Failed, &q.Webhooks.OldestPending); err != nil {
121		return q, err
122	}
123	if err := s.queryEach(`SELECT d.id, `+repoPathExpr+`, w.url, d.attempts, COALESCE(d.last_status, 0),
124		COALESCE(d.last_error, ''), COALESCE(d.failed_at, ''), d.created_at
125		FROM webhook_deliveries d JOIN webhooks w ON w.id = d.webhook_id`+repoJoin("w.repo_id")+`
126		WHERE d.delivered_at IS NULL AND (d.failed_at IS NOT NULL OR d.attempts > 0)
127		ORDER BY d.id DESC LIMIT ?`, func(sc scanner) error {
128		var d QueueDeliveryRow
129		if err := sc.Scan(&d.ID, &d.Repo, &d.URL, &d.Attempts, &d.Status, &d.LastError, &d.FailedAt, &d.CreatedAt); err != nil {
130			return err
131		}
132		q.Webhooks.Items = append(q.Webhooks.Items, d)
133		return nil
134	}); err != nil {
135		return q, err
136	}
137
138	if err := s.DB.QueryRow(`SELECT
139		COUNT(*) FILTER (WHERE sent_at IS NULL AND failed_at IS NULL),
140		COUNT(*) FILTER (WHERE sent_at IS NULL AND failed_at IS NULL AND attempts > 0),
141		COUNT(*) FILTER (WHERE failed_at IS NOT NULL),
142		COALESCE(MIN(created_at) FILTER (WHERE sent_at IS NULL AND failed_at IS NULL), '')
143		FROM notifications`).Scan(&q.Mail.Pending, &q.Mail.Retrying, &q.Mail.Failed, &q.Mail.OldestPending); err != nil {
144		return q, err
145	}
146	if err := s.queryEach(`SELECT id, recipient, subject, attempts, COALESCE(last_error, ''), COALESCE(failed_at, ''), created_at
147		FROM notifications WHERE sent_at IS NULL AND (failed_at IS NOT NULL OR attempts > 0)
148		ORDER BY id DESC LIMIT ?`, func(sc scanner) error {
149		var m QueueMailRow
150		if err := sc.Scan(&m.ID, &m.Recipient, &m.Subject, &m.Attempts, &m.LastError, &m.FailedAt, &m.CreatedAt); err != nil {
151			return err
152		}
153		q.Mail.Items = append(q.Mail.Items, m)
154		return nil
155	}); err != nil {
156		return q, err
157	}
158
159	if err := s.DB.QueryRow(`SELECT COUNT(*) FILTER (WHERE dirty = 1), COUNT(*) FILTER (WHERE last_error != '')
160		FROM mirrors`).Scan(&q.Mirrors.Dirty, &q.Mirrors.Errors); err != nil {
161		return q, err
162	}
163	if err := s.queryEach(`SELECT m.id, `+repoPathExpr+`, m.direction, m.url, m.last_sync, m.last_error
164		FROM mirrors m`+repoJoin("m.repo_id")+` WHERE m.last_error != '' ORDER BY m.id DESC LIMIT ?`, func(sc scanner) error {
165		var m QueueMirrorRow
166		if err := sc.Scan(&m.ID, &m.Repo, &m.Direction, &m.URL, &m.LastSync, &m.LastError); err != nil {
167			return err
168		}
169		q.Mirrors.Items = append(q.Mirrors.Items, m)
170		return nil
171	}); err != nil {
172		return q, err
173	}
174
175	if err := s.DB.QueryRow(`SELECT COUNT(*) FILTER (WHERE status = 'pending'), COUNT(*) FILTER (WHERE status = 'running'),
176		COALESCE(MIN(created_at) FILTER (WHERE status = 'pending'), '') FROM builds`).Scan(&q.Builds.Pending, &q.Builds.Running, &q.Builds.OldestPending); err != nil {
177		return q, err
178	}
179	if err := s.queryEach(`SELECT `+repoPathExpr+`, b.number, b.job, b.started_at
180		FROM builds b`+repoJoin("b.repo_id")+` WHERE b.status = 'running' ORDER BY b.started_at LIMIT ?`, func(sc scanner) error {
181		var b QueueBuildRow
182		if err := sc.Scan(&b.Repo, &b.Number, &b.Job, &b.StartedAt); err != nil {
183			return err
184		}
185		q.Builds.Items = append(q.Builds.Items, b)
186		return nil
187	}); err != nil {
188		return q, err
189	}
190
191	if err := s.DB.QueryRow(`SELECT COUNT(*) FROM dep_checks WHERE last_error != ''`).Scan(&q.Deps.Errors); err != nil {
192		return q, err
193	}
194	if err := s.queryEach(`SELECT `+repoPathExpr+`, c.last_check, c.last_error
195		FROM dep_checks c`+repoJoin("c.repo_id")+` WHERE c.last_error != '' ORDER BY c.repo_id LIMIT ?`, func(sc scanner) error {
196		var d QueueDepRow
197		if err := sc.Scan(&d.Repo, &d.LastCheck, &d.LastError); err != nil {
198			return err
199		}
200		q.Deps.Items = append(q.Deps.Items, d)
201		return nil
202	}); err != nil {
203		return q, err
204	}
205	return q, nil
206}
207
208type scanner interface{ Scan(...any) error }
209
210// queryEach runs a capped item query and hands each row to fn.
211func (s *Store) queryEach(query string, fn func(scanner) error) error {
212	rows, err := s.DB.Query(query, queueItemCap)
213	if err != nil {
214		return err
215	}
216	defer rows.Close()
217	for rows.Next() {
218		if err := fn(rows); err != nil {
219			return err
220		}
221	}
222	return rows.Err()
223}