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}