internal/mailin/mailin.go

460 lines · 14434 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	// 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}