internal/store/notify.go

ada2c022c4701ad3d06e36ff79a54629ed9ff501
gitbay/internal/store/notify.go history · blame · raw

102 lines · 3009 bytes

  1package store
  2
  3import "time"
  4
  5type QueuedMail struct {
  6	ID        int64
  7	Recipient string
  8	Subject   string
  9	Body      string
 10	Attempts  int
 11}
 12
 13func (s *Store) EnqueueMail(recipient, subject, body string) error {
 14	_, err := s.DB.Exec(
 15		"INSERT INTO notifications (recipient, subject, body) VALUES (?, ?, ?)",
 16		recipient, subject, body)
 17	return err
 18}
 19
 20func (s *Store) DueMail(limit int) ([]QueuedMail, error) {
 21	rows, err := s.DB.Query(`
 22		SELECT id, recipient, subject, body, attempts FROM notifications
 23		WHERE sent_at IS NULL AND failed_at IS NULL
 24		  AND (next_attempt_at IS NULL OR next_attempt_at <= ?)
 25		ORDER BY id LIMIT ?`, fmtTime(time.Now()), limit)
 26	if err != nil {
 27		return nil, err
 28	}
 29	defer rows.Close()
 30	var out []QueuedMail
 31	for rows.Next() {
 32		var m QueuedMail
 33		if err := rows.Scan(&m.ID, &m.Recipient, &m.Subject, &m.Body, &m.Attempts); err != nil {
 34			return nil, err
 35		}
 36		out = append(out, m)
 37	}
 38	return out, rows.Err()
 39}
 40
 41func (s *Store) MarkMailSent(id int64) error {
 42	_, err := s.DB.Exec(
 43		"UPDATE notifications SET sent_at = strftime('%Y-%m-%dT%H:%M:%fZ','now'), attempts = attempts + 1 WHERE id = ?", id)
 44	return err
 45}
 46
 47func (s *Store) MarkMailFailed(id int64, errMsg string, nextAt *time.Time) error {
 48	if nextAt == nil {
 49		_, err := s.DB.Exec(
 50			"UPDATE notifications SET failed_at = strftime('%Y-%m-%dT%H:%M:%fZ','now'), attempts = attempts + 1, last_error = ? WHERE id = ?",
 51			errMsg, id)
 52		return err
 53	}
 54	_, err := s.DB.Exec(
 55		"UPDATE notifications SET attempts = attempts + 1, last_error = ?, next_attempt_at = ? WHERE id = ?",
 56		errMsg, fmtTime(*nextAt), id)
 57	return err
 58}
 59
 60// IssueParticipants returns distinct user ids involved in an issue: the
 61// author and every commenter.
 62func (s *Store) IssueParticipants(issueID int64) ([]int64, error) {
 63	return s.idQuery(`
 64		SELECT author_id FROM issues WHERE id = ?
 65		UNION SELECT author_id FROM issue_comments WHERE issue_id = ?`, issueID, issueID)
 66}
 67
 68// MRParticipants returns distinct user ids involved in an MR: author,
 69// commenters, reviewers.
 70func (s *Store) MRParticipants(mrID int64) ([]int64, error) {
 71	return s.idQuery(`
 72		SELECT author_id FROM merge_requests WHERE id = ?
 73		UNION SELECT author_id FROM mr_comments WHERE mr_id = ?
 74		UNION SELECT reviewer_id FROM mr_reviews WHERE mr_id = ?`, mrID, mrID, mrID)
 75}
 76
 77// RepoNotifyTargets returns who should hear about new activity on a repo:
 78// the owning user, or every admin of the owning org.
 79func (s *Store) RepoNotifyTargets(repo Repo) ([]int64, error) {
 80	if repo.OwnerKind == "user" {
 81		return []int64{repo.OwnerID}, nil
 82	}
 83	return s.idQuery(
 84		"SELECT user_id FROM org_members WHERE org_id = ? AND role = 'admin'", repo.OwnerID)
 85}
 86
 87func (s *Store) idQuery(q string, args ...any) ([]int64, error) {
 88	rows, err := s.DB.Query(q, args...)
 89	if err != nil {
 90		return nil, err
 91	}
 92	defer rows.Close()
 93	var out []int64
 94	for rows.Next() {
 95		var id int64
 96		if err := rows.Scan(&id); err != nil {
 97			return nil, err
 98		}
 99		out = append(out, id)
100	}
101	return out, rows.Err()
102}