internal/store/notify.go
126 lines · 4158 bytes
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}