internal/store/webhooks.go
201 lines · 5958 bytes
1package store
2
3import (
4 "fmt"
5 "time"
6)
7
8type Webhook struct {
9 ID int64
10 URL string
11 Secret string
12 Events string // "*" or comma-separated kinds
13 Active bool
14 CreatedAt string
15}
16
17type Delivery struct {
18 ID int64
19 WebhookID int64
20 URL string
21 Secret string
22 EventID int64
23 EventKind string
24 RepoPath string
25 Actor string
26 DataJSON string
27 EventAt string
28 Attempts int
29}
30
31type DeliveryStatus struct {
32 ID int64
33 URL string
34 EventKind string
35 Status string // pending | delivered | failed
36 Attempts int
37 LastStatus int
38 LastError string
39 CreatedAt string
40}
41
42// AddWebhook stores the hook, then seals its secret under the new row's
43// id in the same transaction.
44func (s *Store) AddWebhook(repoID int64, url, secret, events string) (int64, error) {
45 tx, err := s.DB.Begin()
46 if err != nil {
47 return 0, err
48 }
49 defer tx.Rollback()
50 res, err := tx.Exec(
51 "INSERT INTO webhooks (repo_id, url, secret, events) VALUES (?, ?, '', ?)",
52 repoID, url, events)
53 if err != nil {
54 return 0, err
55 }
56 id, err := res.LastInsertId()
57 if err != nil {
58 return 0, err
59 }
60 if secret != "" {
61 sealed, err := s.sealValue(webhookAAD(id), secret)
62 if err != nil {
63 return 0, err
64 }
65 if _, err := tx.Exec("UPDATE webhooks SET secret = ? WHERE id = ?", sealed, id); err != nil {
66 return 0, err
67 }
68 }
69 return id, tx.Commit()
70}
71
72func (s *Store) ListWebhooks(repoID int64) ([]Webhook, error) {
73 rows, err := s.DB.Query(
74 "SELECT id, url, secret, events, active, created_at FROM webhooks WHERE repo_id = ? ORDER BY id", repoID)
75 if err != nil {
76 return nil, err
77 }
78 defer rows.Close()
79 var out []Webhook
80 for rows.Next() {
81 var w Webhook
82 var active int
83 if err := rows.Scan(&w.ID, &w.URL, &w.Secret, &w.Events, &active, &w.CreatedAt); err != nil {
84 return nil, err
85 }
86 if w.Secret, err = s.openValue(webhookAAD(w.ID), w.Secret); err != nil {
87 return nil, fmt.Errorf("webhook %d: %w", w.ID, err)
88 }
89 w.Active = active != 0
90 out = append(out, w)
91 }
92 return out, rows.Err()
93}
94
95func (s *Store) RemoveWebhook(repoID, hookID int64) error {
96 res, err := s.DB.Exec("DELETE FROM webhooks WHERE repo_id = ? AND id = ?", repoID, hookID)
97 if err != nil {
98 return err
99 }
100 if n, _ := res.RowsAffected(); n == 0 {
101 return ErrNotFound
102 }
103 return nil
104}
105
106// DueDeliveries returns pending deliveries whose time has come, with the
107// event and hook context needed to send them.
108func (s *Store) DueDeliveries(limit int) ([]Delivery, error) {
109 rows, err := s.DB.Query(`
110 SELECT d.id, d.webhook_id, w.url, w.secret, d.event_id, e.kind,
111 COALESCE(u2.username, o.name, '') || '/' || COALESCE(r.name, ''),
112 COALESCE(u.username, ''), e.data_json, e.created_at, d.attempts
113 FROM webhook_deliveries d
114 JOIN webhooks w ON w.id = d.webhook_id
115 JOIN events e ON e.id = d.event_id
116 LEFT JOIN users u ON u.id = e.actor_id
117 LEFT JOIN repos r ON r.id = e.repo_id
118 LEFT JOIN users u2 ON r.owner_kind = 'user' AND u2.id = r.owner_id
119 LEFT JOIN orgs o ON r.owner_kind = 'org' AND o.id = r.owner_id
120 WHERE d.delivered_at IS NULL AND d.failed_at IS NULL
121 AND (d.next_attempt_at IS NULL OR d.next_attempt_at <= ?)
122 ORDER BY d.id LIMIT ?`, fmtTime(time.Now()), limit)
123 if err != nil {
124 return nil, err
125 }
126 defer rows.Close()
127 var out []Delivery
128 for rows.Next() {
129 var d Delivery
130 if err := rows.Scan(&d.ID, &d.WebhookID, &d.URL, &d.Secret, &d.EventID, &d.EventKind,
131 &d.RepoPath, &d.Actor, &d.DataJSON, &d.EventAt, &d.Attempts); err != nil {
132 return nil, err
133 }
134 if d.Secret, err = s.openValue(webhookAAD(d.WebhookID), d.Secret); err != nil {
135 return nil, fmt.Errorf("webhook %d: %w", d.WebhookID, err)
136 }
137 out = append(out, d)
138 }
139 return out, rows.Err()
140}
141
142func (s *Store) MarkDelivered(id int64, status int) error {
143 _, err := s.DB.Exec(`
144 UPDATE webhook_deliveries SET delivered_at = strftime('%Y-%m-%dT%H:%M:%fZ','now'),
145 attempts = attempts + 1, last_status = ?, last_error = NULL WHERE id = ?`, status, id)
146 return err
147}
148
149// MarkAttemptFailed records a failed attempt; nextAt nil dead-letters it.
150func (s *Store) MarkAttemptFailed(id int64, status int, errMsg string, nextAt *time.Time) error {
151 if nextAt == nil {
152 _, err := s.DB.Exec(`
153 UPDATE webhook_deliveries SET failed_at = strftime('%Y-%m-%dT%H:%M:%fZ','now'),
154 attempts = attempts + 1, last_status = ?, last_error = ? WHERE id = ?`, status, errMsg, id)
155 return err
156 }
157 _, err := s.DB.Exec(`
158 UPDATE webhook_deliveries SET attempts = attempts + 1, last_status = ?, last_error = ?,
159 next_attempt_at = ? WHERE id = ?`, status, errMsg, fmtTime(*nextAt), id)
160 return err
161}
162
163func (s *Store) ListDeliveries(repoID int64, limit int) ([]DeliveryStatus, error) {
164 rows, err := s.DB.Query(`
165 SELECT d.id, w.url, e.kind,
166 CASE WHEN d.delivered_at IS NOT NULL THEN 'delivered'
167 WHEN d.failed_at IS NOT NULL THEN 'failed'
168 ELSE 'pending' END,
169 d.attempts, COALESCE(d.last_status, 0), COALESCE(d.last_error, ''), d.created_at
170 FROM webhook_deliveries d
171 JOIN webhooks w ON w.id = d.webhook_id
172 JOIN events e ON e.id = d.event_id
173 WHERE w.repo_id = ? ORDER BY d.id DESC LIMIT ?`, repoID, limit)
174 if err != nil {
175 return nil, err
176 }
177 defer rows.Close()
178 var out []DeliveryStatus
179 for rows.Next() {
180 var d DeliveryStatus
181 if err := rows.Scan(&d.ID, &d.URL, &d.EventKind, &d.Status, &d.Attempts, &d.LastStatus, &d.LastError, &d.CreatedAt); err != nil {
182 return nil, err
183 }
184 out = append(out, d)
185 }
186 return out, rows.Err()
187}
188
189// Redeliver resets a delivery for an immediate retry.
190func (s *Store) Redeliver(repoID, deliveryID int64) error {
191 res, err := s.DB.Exec(`
192 UPDATE webhook_deliveries SET delivered_at = NULL, failed_at = NULL, next_attempt_at = NULL
193 WHERE id = ? AND webhook_id IN (SELECT id FROM webhooks WHERE repo_id = ?)`, deliveryID, repoID)
194 if err != nil {
195 return err
196 }
197 if n, _ := res.RowsAffected(); n == 0 {
198 return ErrNotFound
199 }
200 return nil
201}