internal/mailin/mailin.go
392 lines · 12039 bytes
1// Package mailin turns replies to notification mail into comments
2// (#295). A poller reads a mailbox over IMAP; each unseen message
3// addressed to reply+<token>@<domain> is checked (token, account, sender
4// address, the account's access to the thread now) and posted by
5// dispatching issue comment or mr comment as that account. A refusal
6// sends nothing back and is written to the audit log with its reason,
7// never the message's content. Every message is marked seen once it is
8// handled, posted or refused; only a failure that may pass (the
9// database busy) leaves it for the next poll.
10package mailin
11
12import (
13 "bytes"
14 "context"
15 "crypto/sha256"
16 "encoding/hex"
17 "errors"
18 "fmt"
19 "log/slog"
20 "net/mail"
21 "net/textproto"
22 "strconv"
23 "strings"
24 "time"
25
26 "gitbay.org/gitbay/internal/config"
27 "gitbay.org/gitbay/internal/control"
28 "gitbay.org/gitbay/internal/imapc"
29 "gitbay.org/gitbay/internal/mailreply"
30 "gitbay.org/gitbay/internal/protocol"
31 "gitbay.org/gitbay/internal/store"
32)
33
34// Mailbox is what the processor needs of an IMAP session; imapc.Client
35// implements it, and tests use a fake.
36type Mailbox interface {
37 Unseen() ([]uint32, error)
38 Fetch(uid uint32) ([]byte, error)
39 MarkSeen(uid uint32) error
40}
41
42// maxTries is how many polls a message that keeps failing is tried in
43// before it is marked seen and given up on.
44const maxTries = 5
45
46// Processor handles fetched messages.
47type Processor struct {
48 St *store.Store
49 Cfg config.Config
50 Now func() time.Time
51
52 tries map[uint32]int
53 // A window of refusal rows, bounded because anyone can send mail
54 // to the mailbox.
55 windowStart time.Time
56 windowRows int
57}
58
59// refusalsPerMinute bounds the audit rows refusals write.
60const refusalsPerMinute = 60
61
62// Result is what became of one message.
63type Result struct {
64 Posted bool
65 Retry bool // a failure that may pass; the message is left unseen
66 Reason string // why it was refused or failed; empty when posted
67}
68
69func refused(format string, args ...any) Result {
70 return Result{Reason: fmt.Sprintf(format, args...)}
71}
72
73// Drain handles every unseen message in mb. It stops at the first
74// mailbox error.
75func (p *Processor) Drain(mb Mailbox) error {
76 if p.tries == nil {
77 p.tries = map[uint32]int{}
78 }
79 uids, err := mb.Unseen()
80 if err != nil {
81 return err
82 }
83 for _, uid := range uids {
84 // A message that failed in earlier polls before it could be
85 // handled (a fetch the server cut off) is given up on unread.
86 if p.tries[uid] >= maxTries {
87 p.audit(0, "", "gave up after "+strconv.Itoa(maxTries)+" tries")
88 delete(p.tries, uid)
89 if err := mb.MarkSeen(uid); err != nil {
90 return err
91 }
92 continue
93 }
94 raw, err := mb.Fetch(uid)
95 var res Result
96 switch {
97 case errors.Is(err, imapc.ErrTooLarge):
98 res = refused("message larger than %d bytes", imapc.MaxMessage)
99 p.audit(0, "", res.Reason)
100 case errors.Is(err, imapc.ErrLimit):
101 // The server sent more than the limits allow for this
102 // message: it counts a try, and the closed connection ends
103 // the poll.
104 p.tries[uid]++
105 return err
106 case errors.As(err, new(*imapc.RefusedError)):
107 // The server refused this message; the session goes on.
108 res = Result{Retry: true, Reason: "fetch: " + err.Error()}
109 case err != nil:
110 // A connection failure or a timeout says nothing about this
111 // message or the ones after it: the poll ends, no try is
112 // counted.
113 return err
114 default:
115 res = p.Handle(raw)
116 }
117 if res.Retry {
118 p.tries[uid]++
119 if p.tries[uid] < maxTries {
120 slog.Warn("mail reply: will retry", "uid", uid, "err", res.Reason)
121 continue
122 }
123 p.audit(0, "", "gave up after "+strconv.Itoa(maxTries)+" tries: "+res.Reason)
124 }
125 delete(p.tries, uid)
126 if err := mb.MarkSeen(uid); err != nil {
127 return err
128 }
129 }
130 return nil
131}
132
133// Handle checks one message and posts it when every check passes.
134// Refusals are audited here.
135func (p *Processor) Handle(raw []byte) Result {
136 if len(bytes.TrimSpace(raw)) == 0 {
137 return p.refuse(0, "", "empty message")
138 }
139 msg, err := mail.ReadMessage(bytes.NewReader(raw))
140 if err != nil {
141 return p.refuse(0, "", "unreadable message")
142 }
143 msgID := strings.TrimSpace(msg.Header.Get("Message-Id"))
144 if len(msgID) > 200 {
145 msgID = msgID[:200]
146 }
147 if automatic(msg.Header) {
148 return p.refuse(0, msgID, "automatic reply")
149 }
150 in := p.Cfg.Mail.Inbound
151 token := findToken(msg.Header, in.ReplyAddress)
152 if token == "" {
153 return p.refuse(0, msgID, "not addressed to a reply address")
154 }
155 keys := p.St.Keyring()
156 if keys == nil {
157 return Result{Retry: true, Reason: "no secret key loaded"}
158 }
159 secrets, err := keys.Derive(mailreply.Purpose)
160 if err != nil {
161 return Result{Retry: true, Reason: "secret key: " + err.Error()}
162 }
163 target, err := mailreply.Verify(secrets, token, p.now())
164 switch {
165 case errors.Is(err, mailreply.ErrExpired):
166 return p.refuse(target.UserID, msgID, "reply token expired")
167 case err != nil:
168 return p.refuse(0, msgID, err.Error())
169 }
170
171 u, err := p.St.UserByID(target.UserID)
172 switch {
173 case errors.Is(err, store.ErrNotFound):
174 return p.refuse(0, msgID, "account no longer exists")
175 case err != nil:
176 return Result{Retry: true, Reason: err.Error()}
177 case u.Disabled:
178 return p.refuse(u.ID, msgID, "account disabled")
179 case u.Pending:
180 return p.refuse(u.ID, msgID, "account not active")
181 }
182 // Ids are reused after a hard delete: an account created after the
183 // token was minted is not the one it named.
184 if code := p.createdAfter("users", u.ID, target, msgID, "account"); code != nil {
185 return *code
186 }
187 // The token alone is not enough: the reply must come from one of
188 // the account's verified addresses.
189 from, err := msg.Header.AddressList("From")
190 if err != nil || len(from) != 1 {
191 return p.refuse(u.ID, msgID, "no single From address")
192 }
193 ok, err := p.St.VerifiedEmailOf(u.ID, from[0].Address)
194 if err != nil {
195 return Result{Retry: true, Reason: err.Error()}
196 }
197 if !ok {
198 return p.refuse(u.ID, msgID, "From is not a verified address of the account")
199 }
200 if id := p.Cfg.Mail.Inbound.TrustedAuthservID; id != "" {
201 if reason := authenticated(msg.Header, id, from[0].Address); reason != "" {
202 return p.refuse(u.ID, msgID, reason)
203 }
204 }
205 if on, err := p.St.ReplyEnabled(u.ID); err != nil {
206 return Result{Retry: true, Reason: err.Error()}
207 } else if !on {
208 return p.refuse(u.ID, msgID, "reply by mail is off for the account")
209 }
210
211 text, err := textBody(textproto.MIMEHeader(msg.Header), msg.Body, control.MaxCommentBytes)
212 switch {
213 case errors.Is(err, errNoText):
214 return p.refuse(u.ID, msgID, "no text/plain part")
215 case errors.Is(err, errTooLong):
216 return p.refuse(u.ID, msgID, "reply too long")
217 case err != nil:
218 return p.refuse(u.ID, msgID, "unreadable body")
219 }
220 text = stripQuoted(text)
221 if text == "" {
222 return p.refuse(u.ID, msgID, "empty reply")
223 }
224 if len(text) > control.MaxCommentBytes {
225 return p.refuse(u.ID, msgID, "reply too long")
226 }
227
228 repo, err := p.St.RepoByID(target.RepoID)
229 switch {
230 case errors.Is(err, store.ErrNotFound):
231 return p.refuse(u.ID, msgID, "repository no longer exists")
232 case err != nil:
233 return Result{Retry: true, Reason: err.Error()}
234 }
235 if code := p.createdAfter("repos", repo.ID, target, msgID, "repository"); code != nil {
236 return *code
237 }
238
239 // The claim names the thread and the account as well as the
240 // message, so one account's Message-ID cannot suppress another's.
241 id := msgID
242 if id == "" {
243 sum := sha256.Sum256(raw)
244 id = "sha256:" + hex.EncodeToString(sum[:])
245 }
246 key := fmt.Sprintf("%d/%s/%d/%d/%s", u.ID, target.Kind, target.RepoID, target.Number, id)
247 claimed, err := p.St.ClaimMailReply(key)
248 if err != nil {
249 return Result{Retry: true, Reason: err.Error()}
250 }
251 if !claimed {
252 return p.refuse(u.ID, msgID, "already posted")
253 }
254
255 var stdout, stderr bytes.Buffer
256 c := &control.Ctx{User: u, Scope: "full", Store: p.St, Cfg: p.Cfg,
257 Stdin: strings.NewReader(text), Stdout: &stdout, Stderr: &stderr,
258 Source: control.SourceMail}
259 code := control.Dispatch(c, []string{target.Kind, "comment", repo.Path(),
260 strconv.FormatInt(target.Number, 10), "--file", "-"})
261 if code == protocol.ExitOK {
262 return Result{Posted: true}
263 }
264 p.St.ReleaseMailReply(key)
265 reason := strings.TrimSpace(stderr.String())
266 if code == protocol.ExitFailure {
267 return Result{Retry: true, Reason: reason}
268 }
269 // Dispatch has audited a denied or not-found refusal already; this
270 // row says it came by mail and why.
271 return p.refuse(u.ID, msgID, "comment refused: "+reason)
272}
273
274// createdAfter refuses when the row was created after the token was
275// minted: a later account or repository that took a freed id. Created
276// times are compared to the second, the token's precision.
277func (p *Processor) createdAfter(table string, id int64, target mailreply.Target, msgID, what string) *Result {
278 created, err := p.St.CreatedAt(table, id)
279 if err != nil {
280 return &Result{Retry: true, Reason: err.Error()}
281 }
282 if created.Truncate(time.Second).After(target.Issued()) {
283 r := p.refuse(0, msgID, what+" created after the reply token was issued")
284 return &r
285 }
286 return nil
287}
288
289func (p *Processor) now() time.Time {
290 if p.Now != nil {
291 return p.Now()
292 }
293 return time.Now()
294}
295
296func (p *Processor) refuse(actor int64, msgID, reason string) Result {
297 p.audit(actor, msgID, reason)
298 return Result{Reason: reason}
299}
300
301// audit records a refusal: the reason and the Message-ID, never content
302// from the message. Past refusalsPerMinute rows in a minute, the rest of
303// that minute's refusals are counted in one row, written with the first
304// refusal after it.
305func (p *Processor) audit(actor int64, msgID, reason string) {
306 now := p.now()
307 if now.Sub(p.windowStart) >= time.Minute {
308 if over := p.windowRows - refusalsPerMinute; over > 0 {
309 p.St.Audit(0, "refused mail reply", map[string]any{"reason": "throttled", "dropped": over, "source": control.SourceMail})
310 }
311 p.windowStart, p.windowRows = now, 0
312 }
313 p.windowRows++
314 if p.windowRows > refusalsPerMinute {
315 return
316 }
317 data := map[string]any{"reason": reason, "source": control.SourceMail}
318 if msgID != "" {
319 data["message_id"] = msgID
320 }
321 p.St.Audit(actor, "refused mail reply", data)
322}
323
324// automatic reports an auto-responder's message (RFC 3834, and the
325// headers older responders use), which must not post a comment.
326func automatic(h mail.Header) bool {
327 if v := strings.ToLower(strings.TrimSpace(h.Get("Auto-Submitted"))); v != "" && v != "no" {
328 return true
329 }
330 switch strings.ToLower(strings.TrimSpace(h.Get("Precedence"))) {
331 case "bulk", "junk", "list", "auto_reply":
332 return true
333 }
334 return h.Get("X-Autoreply") != "" || h.Get("X-Autorespond") != ""
335}
336
337// findToken returns the reply token from the first recipient header
338// that carries one.
339func findToken(h mail.Header, base string) string {
340 for _, name := range []string{"Delivered-To", "X-Original-To", "Envelope-To", "To", "Cc"} {
341 for _, v := range h[textproto.CanonicalMIMEHeaderKey(name)] {
342 addrs, err := mail.ParseAddressList(v)
343 if err != nil {
344 continue
345 }
346 for _, a := range addrs {
347 if tok, ok := mailreply.TokenFrom(base, a.Address); ok {
348 return tok
349 }
350 }
351 }
352 }
353 return ""
354}
355
356// Poller reads the configured mailbox every poll interval.
357type Poller struct {
358 P *Processor
359 In config.MailInbound
360}
361
362// Run polls until ctx is done.
363func (pl *Poller) Run(ctx context.Context) {
364 t := time.NewTicker(pl.In.Poll())
365 defer t.Stop()
366 for {
367 if err := pl.Once(); err != nil {
368 // The error names the server and the IMAP failure; the
369 // password never reaches it (imapc.Client.Login).
370 slog.Warn("mail reply: poll failed", "server", pl.In.Addr(), "err", err)
371 }
372 select {
373 case <-ctx.Done():
374 return
375 case <-t.C:
376 }
377 }
378}
379
380// Once connects, handles what is waiting, and disconnects.
381func (pl *Poller) Once() error {
382 c, _, err := imapc.Open(pl.In, false, time.Minute)
383 if err != nil {
384 return err
385 }
386 defer c.Close()
387 c.SetDeadline(time.Now().Add(10 * time.Minute))
388 if err := pl.P.Drain(c); err != nil {
389 return err
390 }
391 return pl.P.St.PruneMailReplies(pl.P.now().Add(-mailreply.Lifetime - 24*time.Hour))
392}