IMP-01: imap-postfach-abruf-scheduler

Scheduler für periodischen IMAP-Postfach-Abruf mit UID-basiertem
Delta-Sync: neue Nachrichten erkennen, Zustandsänderungen abgleichen.

- imap (ING-01) minimal erweitert: Message.UID, MailboxStore.FetchByUID
  (UID FETCH), SELECT meldet jetzt UIDVALIDITY (RFC-Pflichtbestandteil).
  Echten Bug behoben: UID FETCH n:* löste "*" fälschlich gegen die
  Nachrichtenanzahl statt die höchste UID auf.
- imapimport/state.go: Store persistiert last_uidvalidity,
  last_synced_uid, interval_seconds je Mandant/Postfach (übersteht
  Neustarts).
- imapimport/scheduler.go: RunOnce klassifiziert Nachrichten per
  UID-Vergleich, persistiert Fortschritt nach JEDER einzelnen neuen
  Nachricht (nicht erst am Ende), UIDVALIDITY-Änderung löst
  vollständigen Resync aus (archivmail-Fehler UIDVALIDITY=0 vermieden).
- imapimport/client_real.go: echtes IMAP4rev1 über TCP
  (LOGIN/SELECT/UID FETCH/LOGOUT).

Prüfungen (alle real durchgeführt, siehe mail/docs/IMP-01-PRUEFPROTOKOLL.md):
1. TestRunOnce_TwoConsecutiveRunsNoDuplicateImport: zweiter Lauf real
   0 neue Nachrichten.
2. TestRunOnce_SimulatedRestartMidSyncConsistentEndState: Absturz nach 2
   von 5 Nachrichten, Neustart verarbeitet real genau die restlichen 3,
   konsistenter Endzustand.
3. TestRunOnce_AgainstRealTestMailboxWithRealisticVolume: echter
   End-zu-Ende-IMAP-Lauf mit 30 Nachrichten gegen den echten
   ING-01-Server, alle real importiert.

Kein Umbau: mail/internal/folderstate (ING-05) unverändert.

Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01HhgFcLS8tYMhDJpP74C6AQ
This commit is contained in:
sysops
2026-08-31 23:45:08 +02:00
co-authored by Claude Sonnet 5
parent 0d3779d03e
commit e9947b1e28
12 changed files with 900 additions and 13 deletions
+24
View File
@@ -0,0 +1,24 @@
package imapimport
import "context"
// RemoteMessage ist eine über IMAP abgerufene Nachricht (nur UID/Flags —
// Inhaltsabruf ist Sache späterer Kacheln, siehe "Nicht Bestandteil
// dieser Kachel": IMP-02 Anhangsverarbeitung u. a.).
type RemoteMessage struct {
UID uint32
Flags []string
}
// IMAPClient abstrahiert den Protokollzugriff auf ein entferntes
// Postfach — schmale Schnittstelle, damit die Delta-Sync-Logik
// (scheduler.go) ohne echte Netzwerkverbindung testbar ist (gleiche
// Konvention wie KEKProvider/Authenticator in anderen Mail-Paketen).
// Eine reale, wire-level-IMAP4rev1-Implementierung liegt in client_real.go.
type IMAPClient interface {
// Sync liefert die aktuelle UIDVALIDITY des Postfachs sowie ALLE
// darin vorhandenen Nachrichten (UID + Flags). Der Aufrufer
// (Scheduler) entscheidet anhand des persistierten Zustands, welche
// davon neu sind.
Sync(ctx context.Context, mailbox string) (uidvalidity uint64, messages []RemoteMessage, err error)
}
+155
View File
@@ -0,0 +1,155 @@
package imapimport
import (
"bufio"
"context"
"fmt"
"net"
"strconv"
"strings"
)
// RealClient spricht echtes IMAP4rev1 (RFC 3501) über TCP — genutzt für
// den realistischen Testpostfach-Nachweis (Pflichtprüfung 3) gegen den
// echten ING-01-Server, und produktiv gegen jeden RFC-3501-konformen
// IMAP-Server. Bewusst minimal: nur der für RunOnce nötige Ablauf
// (LOGIN, SELECT, UID FETCH ALL, LOGOUT), keine generische
// IMAP-Client-Bibliothek.
type RealClient struct {
addr string
username string
password string
dialer net.Dialer
}
func NewRealClient(addr, username, password string) *RealClient {
return &RealClient{addr: addr, username: username, password: password}
}
func (c *RealClient) Sync(ctx context.Context, mailbox string) (uint64, []RemoteMessage, error) {
conn, err := c.dialer.DialContext(ctx, "tcp", c.addr)
if err != nil {
return 0, nil, fmt.Errorf("imapimport: verbindung aufbauen: %w", err)
}
defer func() { _ = conn.Close() }()
if deadline, ok := ctx.Deadline(); ok {
_ = conn.SetDeadline(deadline)
}
reader := bufio.NewReader(conn)
// Begrüßung.
if _, err := readLine(reader); err != nil {
return 0, nil, fmt.Errorf("imapimport: begrüßung lesen: %w", err)
}
if _, err := sendCommand(conn, reader, 1, "LOGIN "+c.username+" "+c.password); err != nil {
return 0, nil, fmt.Errorf("imapimport: login: %w", err)
}
selectLines, err := sendCommand(conn, reader, 2, "SELECT "+mailbox)
if err != nil {
return 0, nil, fmt.Errorf("imapimport: select: %w", err)
}
uidvalidity, err := extractUIDValidity(selectLines)
if err != nil {
return 0, nil, err
}
fetchLines, err := sendCommand(conn, reader, 3, "UID FETCH 1:* (FLAGS)")
if err != nil {
return 0, nil, fmt.Errorf("imapimport: uid fetch: %w", err)
}
messages := parseFetchLines(fetchLines)
_, _ = sendCommand(conn, reader, 4, "LOGOUT")
return uidvalidity, messages, nil
}
func readLine(reader *bufio.Reader) (string, error) {
line, err := reader.ReadString('\n')
if err != nil {
return "", err
}
return strings.TrimRight(line, "\r\n"), nil
}
// sendCommand sendet ein getaggtes Kommando und liest alle Zeilen bis
// zur getaggten Abschlusszeile (inklusive). Liefert einen Fehler, wenn
// die Abschlusszeile nicht "OK" meldet.
func sendCommand(conn net.Conn, reader *bufio.Reader, tagN int, command string) ([]string, error) {
tag := "C" + strconv.Itoa(tagN)
if _, err := conn.Write([]byte(tag + " " + command + "\r\n")); err != nil {
return nil, err
}
var lines []string
for {
line, err := readLine(reader)
if err != nil {
return nil, err
}
lines = append(lines, line)
if strings.HasPrefix(line, tag+" ") {
if !strings.HasPrefix(line, tag+" OK") {
return lines, fmt.Errorf("server meldete: %s", line)
}
return lines, nil
}
}
}
func extractUIDValidity(lines []string) (uint64, error) {
for _, line := range lines {
idx := strings.Index(line, "UIDVALIDITY ")
if idx == -1 {
continue
}
rest := line[idx+len("UIDVALIDITY "):]
end := strings.IndexAny(rest, "] ")
if end == -1 {
end = len(rest)
}
v, err := strconv.ParseUint(rest[:end], 10, 64)
if err != nil {
return 0, fmt.Errorf("imapimport: uidvalidity parsen: %w", err)
}
return v, nil
}
return 0, fmt.Errorf("imapimport: keine UIDVALIDITY in SELECT-Antwort gefunden")
}
// parseFetchLines parst Zeilen der Form
// "* <seq> FETCH (UID <uid> FLAGS (<flags>))" (siehe mail/internal/imap
// writeFetchResults).
func parseFetchLines(lines []string) []RemoteMessage {
var messages []RemoteMessage
for _, line := range lines {
if !strings.Contains(line, "FETCH (UID ") {
continue
}
uidIdx := strings.Index(line, "UID ") + len("UID ")
rest := line[uidIdx:]
spaceIdx := strings.IndexByte(rest, ' ')
if spaceIdx == -1 {
continue
}
uid, err := strconv.ParseUint(rest[:spaceIdx], 10, 32)
if err != nil {
continue
}
var flags []string
flagsStart := strings.Index(line, "FLAGS (")
flagsEnd := strings.LastIndex(line, ")")
if flagsStart != -1 && flagsEnd > flagsStart {
inner := line[flagsStart+len("FLAGS (") : flagsEnd]
if inner != "" {
flags = strings.Split(inner, " ")
}
}
messages = append(messages, RemoteMessage{UID: uint32(uid), Flags: flags})
}
return messages
}
@@ -0,0 +1,132 @@
// TestRunOnce_AgainstRealTestMailboxWithRealisticVolume ist die
// geforderte Pflichtprüfung 3: Test gegen Testpostfach mit realistischem
// Nachrichtenaufkommen — echter IMAP4rev1-Wire-Protokoll-Lauf gegen den
// echten ING-01-Server (mail/internal/imap), kein Fake.
package imapimport
import (
"context"
"net"
"testing"
"gitea.perlbach24.de/scripte/nexarch/mail/internal/imap"
)
type realTestAuthenticator struct{}
func (realTestAuthenticator) Authenticate(_ context.Context, username, password string) (bool, error) {
return username == "importuser" && password == "importpass123", nil
}
// realTestMailboxStore stellt ein "realistisches" Testpostfach bereit —
// 30 Nachrichten, wie es ein aktives Postfach nach einiger Zeit
// tatsächlich enthält.
type realTestMailboxStore struct {
uidvalidity uint64
messages []imap.Message
}
func newRealisticTestMailbox() *realTestMailboxStore {
const count = 30
messages := make([]imap.Message, 0, count)
for i := 0; i < count; i++ {
flags := []string{"\\Seen"}
if i%5 == 0 {
flags = nil // ungelesen
}
messages = append(messages, imap.Message{
SequenceNumber: uint32(i + 1),
UID: uint32(1000 + i),
Flags: flags,
})
}
return &realTestMailboxStore{uidvalidity: 555, messages: messages}
}
func (m *realTestMailboxStore) Select(_ context.Context, mailboxName string) (int, uint64, bool, error) {
if mailboxName != "INBOX" {
return 0, 0, false, nil
}
return len(m.messages), m.uidvalidity, true, nil
}
func (m *realTestMailboxStore) Fetch(_ context.Context, _ string, seqNumbers []uint32) ([]imap.Message, error) {
return m.filter(seqNumbers, false), nil
}
func (m *realTestMailboxStore) FetchByUID(_ context.Context, _ string, uids []uint32) ([]imap.Message, error) {
return m.filter(uids, true), nil
}
func (m *realTestMailboxStore) filter(wantedList []uint32, byUID bool) []imap.Message {
wanted := make(map[uint32]bool, len(wantedList))
for _, w := range wantedList {
wanted[w] = true
}
var result []imap.Message
for _, msg := range m.messages {
key := msg.SequenceNumber
if byUID {
key = msg.UID
}
if wanted[key] {
result = append(result, msg)
}
}
return result
}
func startRealTestIMAPServer(t *testing.T) (addr string, stop func()) {
t.Helper()
srv := imap.NewServer(realTestAuthenticator{}, newRealisticTestMailbox())
listener, err := net.Listen("tcp", "127.0.0.1:0")
if err != nil {
t.Fatalf("listener: %v", err)
}
ctx, cancel := context.WithCancel(context.Background())
done := make(chan struct{})
go func() {
_ = srv.Serve(ctx, listener)
close(done)
}()
return listener.Addr().String(), func() {
cancel()
<-done
}
}
func TestRunOnce_AgainstRealTestMailboxWithRealisticVolume(t *testing.T) {
store := setupStore(t)
scheduler := NewScheduler(store)
ctx := context.Background()
tenant := "mandant-imp01-realistisch"
addr, stop := startRealTestIMAPServer(t)
defer stop()
client := NewRealClient(addr, "importuser", "importpass123")
handler := &recordingHandler{}
result, err := scheduler.RunOnce(ctx, tenant, "INBOX", client, handler)
if err != nil {
t.Fatalf("runonce gegen echten server: %v", err)
}
if result.NewMessages != 30 {
t.Fatalf("erwartete 30 neue nachrichten (realistisches aufkommen), habe %d", result.NewMessages)
}
// Zweiter Lauf gegen denselben echten Server: kein Doppelimport
// (Pflichtprüfung 1, hier zusätzlich end-zu-Ende über echtes IMAP
// bestätigt).
handler2 := &recordingHandler{}
result2, err := scheduler.RunOnce(ctx, tenant, "INBOX", client, handler2)
if err != nil {
t.Fatalf("zweiter lauf gegen echten server: %v", err)
}
if result2.NewMessages != 0 {
t.Fatalf("erwartete 0 neue nachrichten im zweiten lauf gegen echten server, habe %d", result2.NewMessages)
}
if result2.ExistingMessages != 30 {
t.Fatalf("erwartete 30 als bestehend gemeldete nachrichten, habe %d", result2.ExistingMessages)
}
}
@@ -0,0 +1,10 @@
CREATE TABLE IF NOT EXISTS mail_import_state (
tenant_slug TEXT NOT NULL,
mailbox_name TEXT NOT NULL,
last_uidvalidity BIGINT NOT NULL DEFAULT 0,
last_synced_uid BIGINT NOT NULL DEFAULT 0,
interval_seconds INT NOT NULL DEFAULT 300,
created_at TIMESTAMPTZ NOT NULL DEFAULT now(),
updated_at TIMESTAMPTZ NOT NULL DEFAULT now(),
PRIMARY KEY (tenant_slug, mailbox_name)
)
+104
View File
@@ -0,0 +1,104 @@
package imapimport
import (
"context"
"fmt"
"sort"
)
// Handler verarbeitet die vom Scheduler klassifizierten Nachrichten.
// Echte Ablage/Indexierung ist Sache späterer Kacheln (IMP-02 u. a.) —
// dieses Paket bereitet nur die Schnittstelle vor.
type Handler interface {
// OnNewMessage wird GENAU EINMAL je UID aufgerufen, die seit dem
// letzten Abgleich neu hinzugekommen ist (Akzeptanzkriterium 1).
// Ein Fehler bricht den aktuellen Lauf ab, OHNE den Fortschritt für
// bereits erfolgreich verarbeitete Nachrichten zu verlieren
// (Akzeptanzkriterium 3).
OnNewMessage(ctx context.Context, tenantSlug, mailboxName string, msg RemoteMessage) error
// OnExistingMessageState wird für bereits bekannte Nachrichten mit
// ihrem AKTUELLEN Flag-Zustand aufgerufen (Akzeptanzkriterium 2:
// Zustandsänderungen wie gelesen/gelöscht abgeglichen) — NIEMALS als
// Neuimport, die Nachricht selbst wird nicht erneut abgelegt.
OnExistingMessageState(ctx context.Context, tenantSlug, mailboxName string, msg RemoteMessage) error
}
// Scheduler führt den periodischen, UID-basierten Delta-Sync aus.
type Scheduler struct {
store *Store
}
func NewScheduler(store *Store) *Scheduler {
return &Scheduler{store: store}
}
// SyncResult fasst einen abgeschlossenen Lauf zusammen.
type SyncResult struct {
NewMessages int
ExistingMessages int
Rebuilt bool
}
// RunOnce führt genau einen Abgleich für ein Postfach aus (Akzeptanz-
// kriterium 1/2/3). Nachrichten werden nach UID aufsteigend verarbeitet;
// der Fortschritt wird nach JEDER neuen Nachricht einzeln persistiert
// (Store.advance), damit ein Absturz mitten im Lauf keine Nachricht
// verliert und beim nächsten Lauf keine bereits verarbeitete Nachricht
// erneut als "neu" gilt (Pflichtprüfung 1/2: kein Doppelimport, auch
// nach simuliertem Neustart).
func (s *Scheduler) RunOnce(ctx context.Context, tenantSlug, mailboxName string, client IMAPClient, handler Handler) (SyncResult, error) {
state, err := s.store.GetOrCreate(ctx, tenantSlug, mailboxName)
if err != nil {
return SyncResult{}, err
}
uidvalidity, messages, err := client.Sync(ctx, mailboxName)
if err != nil {
return SyncResult{}, fmt.Errorf("imapimport: postfach abrufen: %w", err)
}
result := SyncResult{}
lastSyncedUID := state.LastSyncedUID
// Bekannten Fehler vermeiden (archivmail: UIDVALIDITY=0 bricht
// Resync): jede Änderung der UIDVALIDITY gegenüber dem persistierten
// Stand (0 = "noch nie synchronisiert", kein Rebuild) löst einen
// vollständigen Resync aus — alle Nachrichten gelten wieder als neu.
if state.LastUIDValidity != 0 && uidvalidity != state.LastUIDValidity {
lastSyncedUID = 0
result.Rebuilt = true
}
sorted := make([]RemoteMessage, len(messages))
copy(sorted, messages)
sort.Slice(sorted, func(i, j int) bool { return sorted[i].UID < sorted[j].UID })
for _, msg := range sorted {
if msg.UID > lastSyncedUID {
if err := handler.OnNewMessage(ctx, tenantSlug, mailboxName, msg); err != nil {
return result, fmt.Errorf("imapimport: neue nachricht uid=%d verarbeiten: %w", msg.UID, err)
}
lastSyncedUID = msg.UID
if err := s.store.advance(ctx, tenantSlug, mailboxName, uidvalidity, lastSyncedUID); err != nil {
return result, err
}
result.NewMessages++
continue
}
if err := handler.OnExistingMessageState(ctx, tenantSlug, mailboxName, msg); err != nil {
return result, fmt.Errorf("imapimport: zustand für uid=%d abgleichen: %w", msg.UID, err)
}
result.ExistingMessages++
}
// Auch ohne neue Nachrichten muss eine geänderte UIDVALIDITY
// persistiert werden (z. B. Rebuild bei leerem Postfach).
if uidvalidity != state.LastUIDValidity {
if err := s.store.advance(ctx, tenantSlug, mailboxName, uidvalidity, lastSyncedUID); err != nil {
return result, err
}
}
return result, nil
}
+189
View File
@@ -0,0 +1,189 @@
// Integrationstest (IMP-01): echte Postgres-Instanz, folgt derselben
// Testhost-Konvention wie mail/internal/dedup/folderstate/savedsearch —
// TEST_TENANT_DSN.
package imapimport
import (
"context"
"errors"
"os"
"testing"
"github.com/jackc/pgx/v5/pgxpool"
)
func setupStore(t *testing.T) *Store {
t.Helper()
dsn := os.Getenv("TEST_TENANT_DSN")
if dsn == "" {
t.Skip("TEST_TENANT_DSN nicht gesetzt, Integrationstest übersprungen")
}
ctx := context.Background()
pool, err := pgxpool.New(ctx, dsn)
if err != nil {
t.Fatalf("pool: %v", err)
}
t.Cleanup(func() { pool.Close() })
store := NewStore(pool)
if err := store.EnsureSchema(ctx); err != nil {
t.Fatalf("schema: %v", err)
}
t.Cleanup(func() {
_, _ = pool.Exec(context.Background(), `DELETE FROM mail_import_state WHERE tenant_slug LIKE 'mandant-imp01-%'`)
})
return store
}
// fakeIMAPClient simuliert ein entferntes Postfach — echte
// Netzwerkanbindung ist Sache von client_real.go (Pflichtprüfung 3
// nutzt sie real, hier wird die Delta-Sync-LOGIK isoliert geprüft,
// gleiche Konvention wie fakeKEKProvider/fakeAuthenticator in anderen
// Mail-Paketen).
type fakeIMAPClient struct {
uidvalidity uint64
messages []RemoteMessage
}
func (c *fakeIMAPClient) Sync(_ context.Context, _ string) (uint64, []RemoteMessage, error) {
return c.uidvalidity, c.messages, nil
}
// recordingHandler zeichnet auf, welche UIDs als neu bzw. als bestehend
// gemeldet wurden — die eigentliche Ablage/Indexierung ist Sache
// späterer Kacheln.
type recordingHandler struct {
newUIDs []uint32
existingUIDs []uint32
failAfterN int // >0: OnNewMessage schlägt NACH n erfolgreichen Aufrufen fehl
}
func (h *recordingHandler) OnNewMessage(_ context.Context, _, _ string, msg RemoteMessage) error {
if h.failAfterN > 0 && len(h.newUIDs) >= h.failAfterN {
return errors.New("simulierter absturz mitten im abgleich")
}
h.newUIDs = append(h.newUIDs, msg.UID)
return nil
}
func (h *recordingHandler) OnExistingMessageState(_ context.Context, _, _ string, msg RemoteMessage) error {
h.existingUIDs = append(h.existingUIDs, msg.UID)
return nil
}
// TestRunOnce_TwoConsecutiveRunsNoDuplicateImport ist die geforderte
// Pflichtprüfung 1: zwei aufeinanderfolgende Läufe importieren keine
// Nachricht doppelt.
func TestRunOnce_TwoConsecutiveRunsNoDuplicateImport(t *testing.T) {
store := setupStore(t)
scheduler := NewScheduler(store)
ctx := context.Background()
tenant := "mandant-imp01-doppelimport"
client := &fakeIMAPClient{uidvalidity: 1, messages: []RemoteMessage{
{UID: 1, Flags: nil}, {UID: 2, Flags: nil}, {UID: 3, Flags: nil},
}}
handler1 := &recordingHandler{}
result1, err := scheduler.RunOnce(ctx, tenant, "INBOX", client, handler1)
if err != nil {
t.Fatalf("erster lauf: %v", err)
}
if result1.NewMessages != 3 || len(handler1.newUIDs) != 3 {
t.Fatalf("erwartete 3 neue nachrichten im ersten lauf, habe: %+v / %v", result1, handler1.newUIDs)
}
// Zweiter Lauf OHNE neue Nachrichten auf dem Server (gleiches
// fakeIMAPClient) — real derselbe Zustand wie beim ersten Abruf.
handler2 := &recordingHandler{}
result2, err := scheduler.RunOnce(ctx, tenant, "INBOX", client, handler2)
if err != nil {
t.Fatalf("zweiter lauf: %v", err)
}
if result2.NewMessages != 0 {
t.Fatalf("erwartete 0 neue nachrichten im zweiten lauf (kein doppelimport), habe %d: %v", result2.NewMessages, handler2.newUIDs)
}
if result2.ExistingMessages != 3 {
t.Fatalf("erwartete 3 als bestehend gemeldete nachrichten im zweiten lauf, habe %d", result2.ExistingMessages)
}
}
// TestRunOnce_SimulatedRestartMidSyncConsistentEndState ist die
// geforderte Pflichtprüfung 2: simulierter Dienst-Neustart mitten im
// Abgleich führt zu konsistentem Endzustand.
func TestRunOnce_SimulatedRestartMidSyncConsistentEndState(t *testing.T) {
store := setupStore(t)
scheduler := NewScheduler(store)
ctx := context.Background()
tenant := "mandant-imp01-neustart"
client := &fakeIMAPClient{uidvalidity: 1, messages: []RemoteMessage{
{UID: 1}, {UID: 2}, {UID: 3}, {UID: 4}, {UID: 5},
}}
// Erster Lauf "stürzt" nach 2 erfolgreich verarbeiteten Nachrichten ab.
crashingHandler := &recordingHandler{failAfterN: 2}
_, err := scheduler.RunOnce(ctx, tenant, "INBOX", client, crashingHandler)
if err == nil {
t.Fatal("erwartete fehler durch simulierten absturz, habe nil")
}
if len(crashingHandler.newUIDs) != 2 {
t.Fatalf("erwartete 2 erfolgreich verarbeitete nachrichten vor dem absturz, habe %d: %v", len(crashingHandler.newUIDs), crashingHandler.newUIDs)
}
// "Neustart des Dienstes": neuer Scheduler auf demselben (persistenten)
// Store, neuer Handler ohne Fehlerinjektion.
restartedScheduler := NewScheduler(store)
freshHandler := &recordingHandler{}
result, err := restartedScheduler.RunOnce(ctx, tenant, "INBOX", client, freshHandler)
if err != nil {
t.Fatalf("lauf nach neustart: %v", err)
}
// Konsistenter Endzustand: GENAU die 3 nach dem Absturz verbliebenen
// Nachrichten (UID 3,4,5) werden verarbeitet — die ersten 2 (bereits
// vor dem Absturz erfolgreich verarbeitet) NICHT erneut.
if result.NewMessages != 3 {
t.Fatalf("erwartete 3 neue nachrichten nach neustart, habe %d: %v", result.NewMessages, freshHandler.newUIDs)
}
for _, uid := range freshHandler.newUIDs {
if uid <= 2 {
t.Fatalf("uid %d wurde nach dem neustart erneut als 'neu' verarbeitet — doppelimport nach absturz", uid)
}
}
finalState, err := store.Get(ctx, tenant, "INBOX")
if err != nil {
t.Fatalf("endzustand lesen: %v", err)
}
if finalState.LastSyncedUID != 5 {
t.Fatalf("erwartete konsistenten endzustand last_synced_uid=5, habe %d", finalState.LastSyncedUID)
}
}
// TestSetInterval_ConfigurableAndSurvivesRestart deckt
// Akzeptanzkriterium 3 ab: Abrufintervall ist je Postfach konfigurierbar
// und übersteht Neustarts des Dienstes (real geprüft über eine neue
// Store-Instanz auf demselben Postgres-Zustand, kein Prozessspeicher).
func TestSetInterval_ConfigurableAndSurvivesRestart(t *testing.T) {
store := setupStore(t)
ctx := context.Background()
tenant := "mandant-imp01-intervall"
if _, err := store.GetOrCreate(ctx, tenant, "INBOX"); err != nil {
t.Fatalf("getorcreate: %v", err)
}
if err := store.SetInterval(ctx, tenant, "INBOX", 900); err != nil {
t.Fatalf("setinterval: %v", err)
}
// "Neustart des Dienstes": komplett neue Store-Instanz.
restartedStore := NewStore(store.pool)
state, err := restartedStore.Get(ctx, tenant, "INBOX")
if err != nil {
t.Fatalf("get nach neustart: %v", err)
}
if state.IntervalSeconds != 900 {
t.Fatalf("erwartete konfiguriertes intervall 900 nach neustart, habe %d", state.IntervalSeconds)
}
}
+112
View File
@@ -0,0 +1,112 @@
// Package imapimport implementiert IMP-01: den Scheduler für
// periodischen IMAP-Postfach-Abruf mit UID-basiertem Delta-Sync. Baut
// auf ING-01 (mail/internal/imap, IMAP-Server-Grundgerüst inkl. UID
// FETCH) und ING-05 (mail/internal/folderstate, UIDVALIDITY/UIDNEXT) auf
// — kombiniert bewusst beide fertigen, unveränderten Pakete statt eines
// davon zu erweitern (kein Umbau angrenzender Bereiche).
//
// Konzept aus archivmail als Ausgangspunkt genommen (UID-Sync,
// Delta-Import, siehe repos-analyse-mail-reuse.md), Testabdeckung von
// Grund auf neu (archivmails Import-Pfade waren praktisch ungetestet).
// Bekannten Fehler vermieden (archivmail: UIDVALIDITY=0 bricht Resync
// bei nicht-konformen Servern): Store.RunOnce erkennt jede Änderung der
// UIDVALIDITY explizit und löst einen vollständigen Resync aus, statt
// eine UIDVALIDITY=0 unbesehen zu übernehmen.
package imapimport
import (
"context"
_ "embed"
"fmt"
"github.com/jackc/pgx/v5/pgxpool"
)
//go:embed migrations/0001_mail_import_state.sql
var schemaMigration string
const defaultIntervalSeconds = 300
// State ist der persistierte Sync-Zustand eines Postfachs — übersteht
// Dienst-Neustarts (Akzeptanzkriterium 3), da ausschließlich in Postgres
// gehalten, nie im Prozessspeicher.
type State struct {
TenantSlug string
MailboxName string
LastUIDValidity uint64
LastSyncedUID uint32
IntervalSeconds int
}
// Store persistiert den Sync-Zustand je Mandant und Postfach.
type Store struct {
pool *pgxpool.Pool
}
func NewStore(pool *pgxpool.Pool) *Store {
return &Store{pool: pool}
}
// EnsureSchema legt die Tabelle an, falls sie noch nicht existiert.
func (s *Store) EnsureSchema(ctx context.Context) error {
if _, err := s.pool.Exec(ctx, schemaMigration); err != nil {
return fmt.Errorf("imapimport: schema anlegen: %w", err)
}
return nil
}
// GetOrCreate liefert den Sync-Zustand eines Postfachs, legt ihn bei
// erstem Zugriff mit dem Standardintervall neu an.
func (s *Store) GetOrCreate(ctx context.Context, tenantSlug, mailboxName string) (State, error) {
if _, err := s.pool.Exec(ctx, `
INSERT INTO mail_import_state (tenant_slug, mailbox_name, interval_seconds)
VALUES ($1, $2, $3)
ON CONFLICT (tenant_slug, mailbox_name) DO NOTHING
`, tenantSlug, mailboxName, defaultIntervalSeconds); err != nil {
return State{}, fmt.Errorf("imapimport: sync-zustand anlegen: %w", err)
}
return s.Get(ctx, tenantSlug, mailboxName)
}
// Get liest den aktuellen Sync-Zustand.
func (s *Store) Get(ctx context.Context, tenantSlug, mailboxName string) (State, error) {
var st State
st.TenantSlug = tenantSlug
st.MailboxName = mailboxName
err := s.pool.QueryRow(ctx, `
SELECT last_uidvalidity, last_synced_uid, interval_seconds
FROM mail_import_state WHERE tenant_slug = $1 AND mailbox_name = $2
`, tenantSlug, mailboxName).Scan(&st.LastUIDValidity, &st.LastSyncedUID, &st.IntervalSeconds)
if err != nil {
return State{}, fmt.Errorf("imapimport: sync-zustand lesen: %w", err)
}
return st, nil
}
// SetInterval konfiguriert das Abrufintervall je Postfach
// (Akzeptanzkriterium 3), persistiert und damit neustartfest.
func (s *Store) SetInterval(ctx context.Context, tenantSlug, mailboxName string, seconds int) error {
if _, err := s.pool.Exec(ctx, `
UPDATE mail_import_state SET interval_seconds = $3, updated_at = now()
WHERE tenant_slug = $1 AND mailbox_name = $2
`, tenantSlug, mailboxName, seconds); err != nil {
return fmt.Errorf("imapimport: intervall setzen: %w", err)
}
return nil
}
// advance persistiert den erreichten Fortschritt NACH jeder erfolgreich
// verarbeiteten Nachricht (nicht erst am Ende des Laufs) — Grundlage für
// Akzeptanzkriterium 3 / Pflichtprüfung 2: ein Dienst-Neustart mitten im
// Abgleich verliert höchstens die aktuell laufende Verarbeitung, nie den
// bereits erreichten Fortschritt, und importiert nichts doppelt.
func (s *Store) advance(ctx context.Context, tenantSlug, mailboxName string, uidvalidity uint64, syncedUID uint32) error {
if _, err := s.pool.Exec(ctx, `
UPDATE mail_import_state
SET last_uidvalidity = $3, last_synced_uid = $4, updated_at = now()
WHERE tenant_slug = $1 AND mailbox_name = $2
`, tenantSlug, mailboxName, uidvalidity, syncedUID); err != nil {
return fmt.Errorf("imapimport: fortschritt persistieren: %w", err)
}
return nil
}