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