288 lines
7.5 KiB
Go
288 lines
7.5 KiB
Go
package mailingress
|
|
|
|
import (
|
|
"bytes"
|
|
"context"
|
|
"crypto/tls"
|
|
"encoding/json"
|
|
"errors"
|
|
"fmt"
|
|
"io"
|
|
"net"
|
|
"sort"
|
|
"strings"
|
|
"sync"
|
|
"time"
|
|
|
|
"github.com/emersion/go-imap"
|
|
"github.com/emersion/go-imap/client"
|
|
_ "github.com/emersion/go-message/charset"
|
|
messageMail "github.com/emersion/go-message/mail"
|
|
"github.com/example/notify-gateway/internal/config"
|
|
"github.com/example/notify-gateway/internal/gateway"
|
|
"github.com/example/notify-gateway/internal/model"
|
|
)
|
|
|
|
const maxMailBytes = 1 << 20
|
|
|
|
type checkpoint struct {
|
|
Validity uint32 `json:"validity"`
|
|
UID uint32 `json:"uid"`
|
|
}
|
|
type Status struct {
|
|
ID string `json:"id"`
|
|
At time.Time `json:"at"`
|
|
Error string `json:"error,omitempty"`
|
|
}
|
|
type Poller struct {
|
|
Config *config.Store
|
|
Queue *gateway.Queue
|
|
mu sync.Mutex
|
|
statuses map[string]Status
|
|
}
|
|
|
|
func (p *Poller) Statuses() []Status {
|
|
p.mu.Lock()
|
|
defer p.mu.Unlock()
|
|
out := []Status{}
|
|
for _, s := range p.statuses {
|
|
out = append(out, s)
|
|
}
|
|
return out
|
|
}
|
|
func (p *Poller) Run(ctx context.Context) {
|
|
next := map[string]time.Time{}
|
|
ticker := time.NewTicker(time.Second)
|
|
defer ticker.Stop()
|
|
for {
|
|
for _, m := range p.Config.Get().Ingress.Mail {
|
|
if ctx.Err() != nil {
|
|
return
|
|
}
|
|
if !m.Enabled || time.Now().Before(next[m.ID]) {
|
|
continue
|
|
}
|
|
pollCtx, cancel := context.WithTimeout(ctx, 60*time.Second)
|
|
err := p.Poll(pollCtx, m)
|
|
cancel()
|
|
s := Status{ID: m.ID, At: time.Now().UTC()}
|
|
if err != nil {
|
|
s.Error = "Abruf fehlgeschlagen; Checkpoint bleibt vor der betroffenen Nachricht. Verbindung, MIME-Größe/Text und Zuordnung prüfen."
|
|
}
|
|
p.mu.Lock()
|
|
if p.statuses == nil {
|
|
p.statuses = map[string]Status{}
|
|
}
|
|
p.statuses[m.ID] = s
|
|
p.mu.Unlock()
|
|
next[m.ID] = time.Now().Add(time.Duration(max(m.PollSeconds, 10)) * time.Second)
|
|
}
|
|
select {
|
|
case <-ctx.Done():
|
|
return
|
|
case <-ticker.C:
|
|
}
|
|
}
|
|
}
|
|
|
|
func (p *Poller) Poll(ctx context.Context, m config.MailIngress) error {
|
|
host, _, err := net.SplitHostPort(m.Address)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
return p.poll(ctx, m, &tls.Config{ServerName: host, MinVersion: tls.VersionTLS12})
|
|
}
|
|
|
|
func (p *Poller) poll(ctx context.Context, m config.MailIngress, tlsConfig *tls.Config) error {
|
|
password, err := config.Secret(m.Password)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
raw, err := (&net.Dialer{Timeout: 15 * time.Second}).DialContext(ctx, "tcp", m.Address)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
defer raw.Close()
|
|
stop := context.AfterFunc(ctx, func() { raw.Close() })
|
|
defer stop()
|
|
if err := raw.SetDeadline(time.Now().Add(30 * time.Second)); err != nil {
|
|
return err
|
|
}
|
|
conn := tls.Client(raw, tlsConfig)
|
|
if err := conn.HandshakeContext(ctx); err != nil {
|
|
return err
|
|
}
|
|
c, err := client.New(conn)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
defer c.Terminate()
|
|
c.Timeout = 30 * time.Second
|
|
if err = c.Login(m.Username, password); err != nil {
|
|
return err
|
|
}
|
|
folder := m.Folder
|
|
if folder == "" {
|
|
folder = "INBOX"
|
|
}
|
|
mailbox, err := c.Select(folder, true)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if mailbox.UidValidity == 0 || mailbox.UidNext == 0 {
|
|
return errors.New("server did not provide stable mailbox UIDs")
|
|
}
|
|
cpID := gateway.ScopedKey("imap", m.ID, m.Address, m.Username, folder)
|
|
saved, err := p.Queue.Store.Checkpoint(ctx, cpID)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
cp := checkpoint{}
|
|
if saved != "" {
|
|
if err := json.Unmarshal([]byte(saved), &cp); err != nil {
|
|
return err
|
|
}
|
|
}
|
|
if saved == "" && !m.ImportExisting {
|
|
cp = checkpoint{Validity: mailbox.UidValidity, UID: mailbox.UidNext - 1}
|
|
return p.save(ctx, cpID, cp)
|
|
}
|
|
if cp.Validity != mailbox.UidValidity {
|
|
cp = checkpoint{Validity: mailbox.UidValidity}
|
|
}
|
|
if cp.UID == ^uint32(0) {
|
|
return nil
|
|
}
|
|
criteria := imap.NewSearchCriteria()
|
|
criteria.Uid = new(imap.SeqSet)
|
|
criteria.Uid.AddRange(cp.UID+1, 0)
|
|
ids, err := c.UidSearch(criteria)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
sort.Slice(ids, func(i, j int) bool { return ids[i] < ids[j] })
|
|
processed := 0
|
|
for _, uid := range ids {
|
|
// IMAP n:* may return the largest UID even when it is smaller than n.
|
|
if uid <= cp.UID {
|
|
continue
|
|
}
|
|
if processed >= 100 {
|
|
break
|
|
}
|
|
processed++
|
|
set := new(imap.SeqSet)
|
|
set.AddNum(uid)
|
|
section := &imap.BodySectionName{Peek: true, Partial: []int{0, maxMailBytes + 1}}
|
|
messages := make(chan *imap.Message, 1)
|
|
if err := c.UidFetch(set, []imap.FetchItem{section.FetchItem()}, messages); err != nil {
|
|
return err
|
|
}
|
|
item := <-messages
|
|
if item == nil {
|
|
return errors.New("message disappeared while fetching")
|
|
}
|
|
body := item.GetBody(section)
|
|
if body == nil {
|
|
return errors.New("missing IMAP body")
|
|
}
|
|
b, err := io.ReadAll(io.LimitReader(body, maxMailBytes+1))
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if len(b) > maxMailBytes {
|
|
return errors.New("mail exceeds 1 MiB")
|
|
}
|
|
msg, skip, err := Parse(b, m)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if !skip {
|
|
key := gateway.ScopedKey(cpID, fmt.Sprint(cp.Validity), fmt.Sprint(uid))
|
|
if _, err := p.Queue.Accept(ctx, msg, key); err != nil {
|
|
return err
|
|
}
|
|
}
|
|
cp.UID = uid
|
|
if err := p.save(ctx, cpID, cp); err != nil {
|
|
return err
|
|
}
|
|
}
|
|
return nil
|
|
}
|
|
func (p *Poller) save(ctx context.Context, id string, cp checkpoint) error {
|
|
b, _ := json.Marshal(cp)
|
|
return p.Queue.Store.SetCheckpoint(ctx, id, string(b))
|
|
}
|
|
|
|
// Parse accepts only explicit plain-text MIME bodies; attachments are ignored.
|
|
// Sender headers are routing filters, not proof of sender identity.
|
|
func Parse(data []byte, c config.MailIngress) (model.InboundMessage, bool, error) {
|
|
if len(data) > maxMailBytes {
|
|
return model.InboundMessage{}, false, errors.New("mail exceeds 1 MiB")
|
|
}
|
|
r, err := messageMail.CreateReader(bytes.NewReader(data))
|
|
if err != nil {
|
|
return model.InboundMessage{}, false, err
|
|
}
|
|
defer r.Close()
|
|
if r.Header.Get("X-Notify-Gateway") != "" || r.Header.Get("Auto-Submitted") != "" && !strings.EqualFold(r.Header.Get("Auto-Submitted"), "no") {
|
|
return model.InboundMessage{}, true, nil
|
|
}
|
|
matches := func(field string, allowed []string) bool {
|
|
if len(allowed) == 0 {
|
|
return true
|
|
}
|
|
addresses, err := r.Header.AddressList(field)
|
|
if err != nil {
|
|
return false
|
|
}
|
|
for _, a := range addresses {
|
|
for _, v := range allowed {
|
|
if strings.EqualFold(a.Address, v) {
|
|
return true
|
|
}
|
|
}
|
|
}
|
|
return false
|
|
}
|
|
if !matches("From", c.From) || !(len(c.To) == 0 || matches("To", c.To) || matches("Cc", c.To)) {
|
|
return model.InboundMessage{}, true, nil
|
|
}
|
|
subject, err := r.Header.Subject()
|
|
if err != nil {
|
|
return model.InboundMessage{}, false, err
|
|
}
|
|
var text strings.Builder
|
|
for {
|
|
part, err := r.NextPart()
|
|
if err == io.EOF {
|
|
break
|
|
}
|
|
if err != nil {
|
|
return model.InboundMessage{}, false, err
|
|
}
|
|
if h, ok := part.Header.(*messageMail.InlineHeader); ok {
|
|
typ, _, err := h.ContentType()
|
|
if err != nil {
|
|
return model.InboundMessage{}, false, err
|
|
}
|
|
if typ == "text/plain" {
|
|
b, err := io.ReadAll(io.LimitReader(part.Body, maxMailBytes+1))
|
|
if err != nil {
|
|
return model.InboundMessage{}, false, err
|
|
}
|
|
if text.Len()+len(b) > maxMailBytes {
|
|
return model.InboundMessage{}, false, errors.New("decoded mail too large")
|
|
}
|
|
text.Write(b)
|
|
text.WriteByte('\n')
|
|
}
|
|
}
|
|
}
|
|
if strings.TrimSpace(text.String()) == "" {
|
|
return model.InboundMessage{}, false, errors.New("mail has no plain-text body")
|
|
}
|
|
return model.InboundMessage{Source: "mail", Channel: c.Channel, Title: subject, Message: strings.TrimSpace(text.String()), Raw: map[string]any{"from": r.Header.Get("From"), "to": r.Header.Get("To"), "message_id": r.Header.Get("Message-ID")}, ReceivedAt: time.Now().UTC()}, false, nil
|
|
}
|