internal/store/queues.go
269 lines · 10184 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 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}