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}