Files
2026-09-16 06:26:16 +02:00

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
}