internal/store/queues.go
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}