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
105 lines
3.8 KiB
Go
105 lines
3.8 KiB
Go
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
|
|
}
|