internal/mailin/mailin.go
460 lines · 14434 bytes
19 symbols in this file
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 // LookupTXT resolves DKIM selector keys; nil is the system
52 // resolver.
53 LookupTXT LookupTXT
54
55 keys keyCache
56 tries map[uint32]int
57 // A window of refusal rows, bounded because anyone can send mail
58 // to the mailbox.
59 windowStart time.Time
60 windowRows int
61}
62
63// refusalsPerMinute bounds the audit rows refusals write.
64const refusalsPerMinute = 60
65
66// Result is what became of one message.
67type Result struct {
68 Posted bool
69 Retry bool // a failure that may pass; the message is left unseen
70 Reason string // why it was refused or failed; empty when posted
71}
72
73func refused(format string, args ...any) Result {
74 return Result{Reason: fmt.Sprintf(format, args...)}
75}
76
77// Drain handles every unseen message in mb. It stops at the first
78// mailbox error.
79func (p *Processor) Drain(mb Mailbox) error {
80 if p.tries == nil {
81 p.tries = map[uint32]int{}
82 }
83 uids, err := mb.Unseen()
84 if err != nil {
85 return err
86 }
87 for _, uid := range uids {
88 // A message that failed in earlier polls before it could be
89 // handled (a fetch the server cut off) is given up on unread.
90 if p.tries[uid] >= maxTries {
91 p.audit(0, "", "gave up after "+strconv.Itoa(maxTries)+" tries")
92 delete(p.tries, uid)
93 if err := mb.MarkSeen(uid); err != nil {
94 return err
95 }
96 continue
97 }
98 raw, err := mb.Fetch(uid)
99 var res Result
100 switch {
101 case errors.Is(err, imapc.ErrTooLarge):
102 res = refused("message larger than %d bytes", imapc.MaxMessage)
103 p.audit(0, "", res.Reason)
104 case errors.Is(err, imapc.ErrLimit):
105 // The server sent more than the limits allow for this
106 // message: it counts a try, and the closed connection ends
107 // the poll.
108 p.tries[uid]++
109 return err
110 case errors.As(err, new(*imapc.RefusedError)):
111 // The server refused this message; the session goes on.
112 res = Result{Retry: true, Reason: "fetch: " + err.Error()}
113 case err != nil:
114 // A connection failure or a timeout says nothing about this
115 // message or the ones after it: the poll ends, no try is
116 // counted.
117 return err
118 default:
119 res = p.Handle(raw)
120 }
121 if res.Retry {
122 p.tries[uid]++
123 if p.tries[uid] < maxTries {
124 slog.Warn("mail reply: will retry", "uid", uid, "err", res.Reason)
125 continue
126 }
127 p.audit(0, "", "gave up after "+strconv.Itoa(maxTries)+" tries: "+res.Reason)
128 }
129 delete(p.tries, uid)
130 if err := mb.MarkSeen(uid); err != nil {
131 return err
132 }
133 }
134 return nil
135}
136
137// Handle checks one message and posts it when every check passes.
138// Refusals are audited here.
139func (p *Processor) Handle(raw []byte) Result {
140 if len(bytes.TrimSpace(raw)) == 0 {
141 return p.refuse(0, "", "empty message")
142 }
143 rh, reason := parseRawHeader(raw)
144 if reason != "" {
145 return p.refuse(0, "", reason)
146 }
147 msg, err := mail.ReadMessage(bytes.NewReader(raw))
148 if err != nil {
149 return p.refuse(0, "", "unreadable message")
150 }
151 msgID := strings.TrimSpace(msg.Header.Get("Message-Id"))
152 if len(msgID) > 200 {
153 msgID = msgID[:200]
154 }
155 if automatic(msg.Header) {
156 return p.refuse(0, msgID, "automatic reply")
157 }
158 in := p.Cfg.Mail.Inbound
159 token, _ := findToken(msg.Header, in.ReplyAddress, recipientFields)
160 if token == "" {
161 return p.refuse(0, msgID, "not addressed to a reply address")
162 }
163 keys := p.St.Keyring()
164 if keys == nil {
165 return Result{Retry: true, Reason: "no secret key loaded"}
166 }
167 secrets, err := keys.Derive(mailreply.Purpose)
168 if err != nil {
169 return Result{Retry: true, Reason: "secret key: " + err.Error()}
170 }
171 target, err := mailreply.Verify(secrets, token, p.now())
172 switch {
173 case errors.Is(err, mailreply.ErrExpired):
174 return p.refuse(target.UserID, msgID, "reply token expired")
175 case err != nil:
176 return p.refuse(0, msgID, err.Error())
177 }
178
179 u, err := p.St.UserByID(target.UserID)
180 switch {
181 case errors.Is(err, store.ErrNotFound):
182 return p.refuse(0, msgID, "account no longer exists")
183 case err != nil:
184 return Result{Retry: true, Reason: err.Error()}
185 case u.Disabled:
186 return p.refuse(u.ID, msgID, "account disabled")
187 case u.Pending:
188 return p.refuse(u.ID, msgID, "account not active")
189 }
190 // An id freed before ids stopped being reused (#306) may have been
191 // taken: an account created after the token was minted is not the
192 // one it named.
193 if code := p.createdAfter("users", u.ID, target, msgID, "account"); code != nil {
194 return *code
195 }
196 // The token alone is not enough: the reply must come from one of
197 // the account's verified addresses.
198 from, err := msg.Header.AddressList("From")
199 if err != nil || len(from) != 1 {
200 return p.refuse(u.ID, msgID, "no single From address")
201 }
202 ok, err := p.St.VerifiedEmailOf(u.ID, from[0].Address)
203 if err != nil {
204 return Result{Retry: true, Reason: err.Error()}
205 }
206 if !ok {
207 return p.refuse(u.ID, msgID, "From is not a verified address of the account")
208 }
209 sigIDs, res := p.authenticate(raw, rh, msg.Header, from[0].Address, token)
210 if res != nil {
211 if res.Retry {
212 return *res
213 }
214 return p.refuse(u.ID, msgID, res.Reason)
215 }
216 if on, err := p.St.ReplyEnabled(u.ID); err != nil {
217 return Result{Retry: true, Reason: err.Error()}
218 } else if !on {
219 return p.refuse(u.ID, msgID, "reply by mail is off for the account")
220 }
221
222 text, err := textBody(textproto.MIMEHeader(msg.Header), msg.Body, control.MaxCommentBytes)
223 switch {
224 case errors.Is(err, errNoText):
225 return p.refuse(u.ID, msgID, "no text/plain part")
226 case errors.Is(err, errTooLong):
227 return p.refuse(u.ID, msgID, "reply too long")
228 case err != nil:
229 return p.refuse(u.ID, msgID, "unreadable body")
230 }
231 text = stripQuoted(text)
232 if text == "" {
233 return p.refuse(u.ID, msgID, "empty reply")
234 }
235 if len(text) > control.MaxCommentBytes {
236 return p.refuse(u.ID, msgID, "reply too long")
237 }
238
239 repo, err := p.St.RepoByID(target.RepoID)
240 switch {
241 case errors.Is(err, store.ErrNotFound):
242 return p.refuse(u.ID, msgID, "repository no longer exists")
243 case err != nil:
244 return Result{Retry: true, Reason: err.Error()}
245 }
246 if code := p.createdAfter("repos", repo.ID, target, msgID, "repository"); code != nil {
247 return *code
248 }
249
250 // The claim names the thread and the account as well as the
251 // message, so one account's Message-ID cannot suppress another's.
252 id := msgID
253 if id == "" {
254 sum := sha256.Sum256(raw)
255 id = "sha256:" + hex.EncodeToString(sum[:])
256 }
257 // Each passing DKIM signature is claimed too, so a copy of a signed
258 // message is not posted again under another Message-ID or with
259 // unsigned fields changed.
260 var claims []string
261 for _, id := range append([]string{id}, sigIDs...) {
262 claims = append(claims, fmt.Sprintf("%d/%s/%d/%d/%s", u.ID, target.Kind, target.RepoID, target.Number, id))
263 }
264 for n, k := range claims {
265 claimed, err := p.St.ClaimMailReply(k)
266 if err != nil || !claimed {
267 for _, k := range claims[:n] {
268 p.St.ReleaseMailReply(k)
269 }
270 if err != nil {
271 return Result{Retry: true, Reason: err.Error()}
272 }
273 return p.refuse(u.ID, msgID, "already posted")
274 }
275 }
276
277 var stdout, stderr bytes.Buffer
278 c := &control.Ctx{User: u, Scope: "full", Store: p.St, Cfg: p.Cfg,
279 Stdin: strings.NewReader(text), Stdout: &stdout, Stderr: &stderr,
280 Source: control.SourceMail}
281 code := control.Dispatch(c, []string{target.Kind, "comment", repo.Path(),
282 strconv.FormatInt(target.Number, 10), "--file", "-"})
283 if code == protocol.ExitOK {
284 return Result{Posted: true}
285 }
286 for _, k := range claims {
287 p.St.ReleaseMailReply(k)
288 }
289 reason = strings.TrimSpace(stderr.String())
290 if code == protocol.ExitFailure {
291 return Result{Retry: true, Reason: reason}
292 }
293 // Dispatch has audited a denied or not-found refusal already; this
294 // row says it came by mail and why.
295 return p.refuse(u.ID, msgID, "comment refused: "+reason)
296}
297
298// authenticate checks that the mail host or the sender's domain vouches
299// for From: an Authentication-Results pass from trusted_authserv_id, or
300// a DKIM signature that verifies here with require_dkim. When both are
301// set either is enough. A DKIM pass also needs the reply address in a
302// signed To or Cc; an Authentication-Results pass takes it from any
303// recipient field. It returns nil when From is authenticated or neither
304// is set, and the unaudited refusal or retry otherwise; sigIDs are the
305// passing DKIM signatures when DKIM authenticated the reply.
306func (p *Processor) authenticate(raw []byte, rh rawHeader, h mail.Header, from, token string) (sigIDs []string, res *Result) {
307 in := p.Cfg.Mail.Inbound
308 var reasons []string
309 if id := in.TrustedAuthservID; id != "" {
310 reason := authenticated(h, id, from)
311 if reason == "" {
312 return nil, nil
313 }
314 reasons = append(reasons, reason)
315 }
316 if in.RequireDKIM {
317 // The reply address must be in a field the signature covers,
318 // or a signed message could be redirected to any token.
319 tok, field := findToken(h, in.ReplyAddress, []string{"To", "Cc"})
320 if tok != token {
321 reasons = append(reasons, "DKIM: reply address not in To or Cc")
322 } else {
323 ids, reason, retry := p.dkimVerified(raw, rh, from, field)
324 if reason == "" {
325 return ids, nil
326 }
327 if retry {
328 return nil, &Result{Retry: true, Reason: reason}
329 }
330 reasons = append(reasons, reason)
331 }
332 }
333 if len(reasons) == 0 {
334 return nil, nil
335 }
336 return nil, &Result{Reason: strings.Join(reasons, "; ")}
337}
338
339// createdAfter refuses when the row was created after the token was
340// minted: a later account or repository that took a freed id. Created
341// times are compared to the second, the token's precision.
342func (p *Processor) createdAfter(table string, id int64, target mailreply.Target, msgID, what string) *Result {
343 created, err := p.St.CreatedAt(table, id)
344 if err != nil {
345 return &Result{Retry: true, Reason: err.Error()}
346 }
347 if created.Truncate(time.Second).After(target.Issued()) {
348 r := p.refuse(0, msgID, what+" created after the reply token was issued")
349 return &r
350 }
351 return nil
352}
353
354func (p *Processor) now() time.Time {
355 if p.Now != nil {
356 return p.Now()
357 }
358 return time.Now()
359}
360
361func (p *Processor) refuse(actor int64, msgID, reason string) Result {
362 p.audit(actor, msgID, reason)
363 return Result{Reason: reason}
364}
365
366// audit records a refusal: the reason and the Message-ID, never content
367// from the message. Past refusalsPerMinute rows in a minute, the rest of
368// that minute's refusals are counted in one row, written with the first
369// refusal after it.
370func (p *Processor) audit(actor int64, msgID, reason string) {
371 now := p.now()
372 if now.Sub(p.windowStart) >= time.Minute {
373 if over := p.windowRows - refusalsPerMinute; over > 0 {
374 p.St.Audit(0, "refused mail reply", map[string]any{"reason": "throttled", "dropped": over, "source": control.SourceMail})
375 }
376 p.windowStart, p.windowRows = now, 0
377 }
378 p.windowRows++
379 if p.windowRows > refusalsPerMinute {
380 return
381 }
382 data := map[string]any{"reason": reason, "source": control.SourceMail}
383 if msgID != "" {
384 data["message_id"] = msgID
385 }
386 p.St.Audit(actor, "refused mail reply", data)
387}
388
389// automatic reports an auto-responder's message (RFC 3834, and the
390// headers older responders use), which must not post a comment.
391func automatic(h mail.Header) bool {
392 if v := strings.ToLower(strings.TrimSpace(h.Get("Auto-Submitted"))); v != "" && v != "no" {
393 return true
394 }
395 switch strings.ToLower(strings.TrimSpace(h.Get("Precedence"))) {
396 case "bulk", "junk", "list", "auto_reply":
397 return true
398 }
399 return h.Get("X-Autoreply") != "" || h.Get("X-Autorespond") != ""
400}
401
402// recipientFields are where a reply address is looked for, in order.
403var recipientFields = []string{"Delivered-To", "X-Original-To", "Envelope-To", "To", "Cc"}
404
405// findToken returns the reply token from the first of names that
406// carries one, and that field's name in lower case.
407func findToken(h mail.Header, base string, names []string) (token, field string) {
408 for _, name := range names {
409 for _, v := range h[textproto.CanonicalMIMEHeaderKey(name)] {
410 addrs, err := mail.ParseAddressList(v)
411 if err != nil {
412 continue
413 }
414 for _, a := range addrs {
415 if tok, ok := mailreply.TokenFrom(base, a.Address); ok {
416 return tok, strings.ToLower(name)
417 }
418 }
419 }
420 }
421 return "", ""
422}
423
424// Poller reads the configured mailbox every poll interval.
425type Poller struct {
426 P *Processor
427 In config.MailInbound
428}
429
430// Run polls until ctx is done.
431func (pl *Poller) Run(ctx context.Context) {
432 t := time.NewTicker(pl.In.Poll())
433 defer t.Stop()
434 for {
435 if err := pl.Once(); err != nil {
436 // The error names the server and the IMAP failure; the
437 // password never reaches it (imapc.Client.Login).
438 slog.Warn("mail reply: poll failed", "server", pl.In.Addr(), "err", err)
439 }
440 select {
441 case <-ctx.Done():
442 return
443 case <-t.C:
444 }
445 }
446}
447
448// Once connects, handles what is waiting, and disconnects.
449func (pl *Poller) Once() error {
450 c, _, err := imapc.Open(pl.In, false, time.Minute)
451 if err != nil {
452 return err
453 }
454 defer c.Close()
455 c.SetDeadline(time.Now().Add(10 * time.Minute))
456 if err := pl.P.Drain(c); err != nil {
457 return err
458 }
459 return pl.P.St.PruneMailReplies(pl.P.now().Add(-mailreply.Lifetime - 24*time.Hour))
460}