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
190 lines
6.6 KiB
Go
190 lines
6.6 KiB
Go
// 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)
|
|
}
|
|
}
|