internal/store/notify.go
118 lines · 3748 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, every commenter, and everyone mentioned.
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 = ?
66 UNION SELECT user_id FROM mentions WHERE kind = 'issue' AND item_id = ?`, issueID, issueID, issueID)
67}
68
69// MRParticipants returns distinct user ids involved in an MR: author,
70// commenters, reviewers, and everyone mentioned.
71func (s *Store) MRParticipants(mrID int64) ([]int64, error) {
72 return s.idQuery(`
73 SELECT author_id FROM merge_requests WHERE id = ?
74 UNION SELECT author_id FROM mr_comments WHERE mr_id = ?
75 UNION SELECT reviewer_id FROM mr_reviews WHERE mr_id = ?
76 UNION SELECT author_id FROM mr_diff_comments WHERE mr_id = ?
77 UNION SELECT user_id FROM mentions WHERE kind = 'mr' AND item_id = ?`, mrID, mrID, mrID, mrID, mrID)
78}
79
80// AddMentions records accounts mentioned in an issue or merge request
81// (kind "issue" or "mr"); a repeat mention is not an error.
82func (s *Store) AddMentions(repoID int64, kind string, itemID int64, userIDs []int64) error {
83 for _, id := range userIDs {
84 if _, err := s.DB.Exec(
85 "INSERT INTO mentions (repo_id, kind, item_id, user_id) VALUES (?, ?, ?, ?) ON CONFLICT DO NOTHING",
86 repoID, kind, itemID, id); err != nil {
87 return err
88 }
89 }
90 return nil
91}
92
93// RepoNotifyTargets returns who should hear about new activity on a repo:
94// the owning user, or every admin of the owning org.
95func (s *Store) RepoNotifyTargets(repo Repo) ([]int64, error) {
96 if repo.OwnerKind == "user" {
97 return []int64{repo.OwnerID}, nil
98 }
99 return s.idQuery(
100 "SELECT user_id FROM org_members WHERE org_id = ? AND role = 'admin'", repo.OwnerID)
101}
102
103func (s *Store) idQuery(q string, args ...any) ([]int64, error) {
104 rows, err := s.DB.Query(q, args...)
105 if err != nil {
106 return nil, err
107 }
108 defer rows.Close()
109 var out []int64
110 for rows.Next() {
111 var id int64
112 if err := rows.Scan(&id); err != nil {
113 return nil, err
114 }
115 out = append(out, id)
116 }
117 return out, rows.Err()
118}