internal/store/notify.go
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}