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