internal/store/notify.go

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

126 lines · 4158 bytes

11 symbols in this file
  1package store
  2
  3import "time"
  4
  5type QueuedMail struct {
  6	ID        int64
  7	Recipient string
  8	ReplyTo   string // "" for none
  9	Subject   string
 10	Body      string
 11	Attempts  int
 12}
 13
 14func (s *Store) EnqueueMail(recipient, subject, body string) error {
 15	return s.EnqueueMailReplyTo(recipient, "", subject, body)
 16}
 17
 18// EnqueueMailReplyTo queues mail with a Reply-To address. The address
 19// carries a reply token, so it is blanked once the row is sent or
 20// dead-lettered.
 21func (s *Store) EnqueueMailReplyTo(recipient, replyTo, subject, body string) error {
 22	_, err := s.DB.Exec(
 23		"INSERT INTO notifications (recipient, reply_to, subject, body) VALUES (?, ?, ?, ?)",
 24		recipient, replyTo, subject, body)
 25	return err
 26}
 27
 28func (s *Store) DueMail(limit int) ([]QueuedMail, error) {
 29	rows, err := s.DB.Query(`
 30		SELECT id, recipient, reply_to, subject, body, attempts FROM notifications
 31		WHERE sent_at IS NULL AND failed_at IS NULL
 32		  AND (next_attempt_at IS NULL OR next_attempt_at <= ?)
 33		ORDER BY id LIMIT ?`, fmtTime(time.Now()), limit)
 34	if err != nil {
 35		return nil, err
 36	}
 37	defer rows.Close()
 38	var out []QueuedMail
 39	for rows.Next() {
 40		var m QueuedMail
 41		if err := rows.Scan(&m.ID, &m.Recipient, &m.ReplyTo, &m.Subject, &m.Body, &m.Attempts); err != nil {
 42			return nil, err
 43		}
 44		out = append(out, m)
 45	}
 46	return out, rows.Err()
 47}
 48
 49func (s *Store) MarkMailSent(id int64) error {
 50	_, err := s.DB.Exec(
 51		"UPDATE notifications SET sent_at = strftime('%Y-%m-%dT%H:%M:%fZ','now'), attempts = attempts + 1, reply_to = '' WHERE id = ?", id)
 52	return err
 53}
 54
 55func (s *Store) MarkMailFailed(id int64, errMsg string, nextAt *time.Time) error {
 56	if nextAt == nil {
 57		_, err := s.DB.Exec(
 58			"UPDATE notifications SET failed_at = strftime('%Y-%m-%dT%H:%M:%fZ','now'), attempts = attempts + 1, last_error = ?, reply_to = '' WHERE id = ?",
 59			errMsg, id)
 60		return err
 61	}
 62	_, err := s.DB.Exec(
 63		"UPDATE notifications SET attempts = attempts + 1, last_error = ?, next_attempt_at = ? WHERE id = ?",
 64		errMsg, fmtTime(*nextAt), id)
 65	return err
 66}
 67
 68// IssueParticipants returns distinct user ids involved in an issue: the
 69// author, every commenter, and everyone mentioned.
 70func (s *Store) IssueParticipants(issueID int64) ([]int64, error) {
 71	return s.idQuery(`
 72		SELECT author_id FROM issues WHERE id = ?
 73		UNION SELECT author_id FROM issue_comments WHERE issue_id = ?
 74		UNION SELECT user_id FROM mentions WHERE kind = 'issue' AND item_id = ?`, issueID, issueID, issueID)
 75}
 76
 77// MRParticipants returns distinct user ids involved in an MR: author,
 78// commenters, reviewers, and everyone mentioned.
 79func (s *Store) MRParticipants(mrID int64) ([]int64, error) {
 80	return s.idQuery(`
 81		SELECT author_id FROM merge_requests WHERE id = ?
 82		UNION SELECT author_id FROM mr_comments WHERE mr_id = ?
 83		UNION SELECT reviewer_id FROM mr_reviews WHERE mr_id = ?
 84		UNION SELECT author_id FROM mr_diff_comments WHERE mr_id = ?
 85		UNION SELECT user_id FROM mentions WHERE kind = 'mr' AND item_id = ?`, mrID, mrID, mrID, mrID, mrID)
 86}
 87
 88// AddMentions records accounts mentioned in an issue or merge request
 89// (kind "issue" or "mr"); a repeat mention is not an error.
 90func (s *Store) AddMentions(repoID int64, kind string, itemID int64, userIDs []int64) error {
 91	for _, id := range userIDs {
 92		if _, err := s.DB.Exec(
 93			"INSERT INTO mentions (repo_id, kind, item_id, user_id) VALUES (?, ?, ?, ?) ON CONFLICT DO NOTHING",
 94			repoID, kind, itemID, id); err != nil {
 95			return err
 96		}
 97	}
 98	return nil
 99}
100
101// RepoNotifyTargets returns who should hear about new activity on a repo:
102// the owning user, or every admin of the owning org.
103func (s *Store) RepoNotifyTargets(repo Repo) ([]int64, error) {
104	if repo.OwnerKind == "user" {
105		return []int64{repo.OwnerID}, nil
106	}
107	return s.idQuery(
108		"SELECT user_id FROM org_members WHERE org_id = ? AND role = 'admin'", repo.OwnerID)
109}
110
111func (s *Store) idQuery(q string, args ...any) ([]int64, error) {
112	rows, err := s.DB.Query(q, args...)
113	if err != nil {
114		return nil, err
115	}
116	defer rows.Close()
117	var out []int64
118	for rows.Next() {
119		var id int64
120		if err := rows.Scan(&id); err != nil {
121			return nil, err
122		}
123		out = append(out, id)
124	}
125	return out, rows.Err()
126}