internal/store/notify.go

5ffa892991a5437dc53680ea1609fabee4b85d42
gitbay/internal/store/notify.go history · blame · raw

103 lines · 3078 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 = ?
 75		UNION SELECT author_id FROM mr_diff_comments WHERE mr_id = ?`, mrID, mrID, mrID, mrID)
 76}
 77
 78// RepoNotifyTargets returns who should hear about new activity on a repo:
 79// the owning user, or every admin of the owning org.
 80func (s *Store) RepoNotifyTargets(repo Repo) ([]int64, error) {
 81	if repo.OwnerKind == "user" {
 82		return []int64{repo.OwnerID}, nil
 83	}
 84	return s.idQuery(
 85		"SELECT user_id FROM org_members WHERE org_id = ? AND role = 'admin'", repo.OwnerID)
 86}
 87
 88func (s *Store) idQuery(q string, args ...any) ([]int64, error) {
 89	rows, err := s.DB.Query(q, args...)
 90	if err != nil {
 91		return nil, err
 92	}
 93	defer rows.Close()
 94	var out []int64
 95	for rows.Next() {
 96		var id int64
 97		if err := rows.Scan(&id); err != nil {
 98			return nil, err
 99		}
100		out = append(out, id)
101	}
102	return out, rows.Err()
103}