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
113 lines
4.3 KiB
Go
113 lines
4.3 KiB
Go
// 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
|
|
}
|