internal/store/webhooks.go

8dcfa45a8ac03a5ff9c36828274d05acadcf846c
gitbay/internal/store/webhooks.go history · blame · raw

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}