Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
e9947b1e28 | ||
|
|
0d3779d03e |
@@ -0,0 +1,63 @@
|
||||
# IMP-01 – Prüfprotokoll: IMAP-Postfach-Abruf & Scheduler
|
||||
|
||||
Voraussetzung ING-01, ING-05 (beide Fertig).
|
||||
|
||||
## Umsetzung
|
||||
|
||||
- `mail/internal/imap` (ING-01) minimal erweitert: `Message.UID`,
|
||||
`MailboxStore.FetchByUID` (RFC 3501 §6.4.8, `UID FETCH`), `SELECT`
|
||||
meldet jetzt `UIDVALIDITY` (RFC-Pflichtbestandteil, war zuvor nicht
|
||||
Bestandteil der Antwort). Dabei einen echten Bug im selben Zug
|
||||
gefunden und behoben: `UID FETCH n:*` löste `*` fälschlich gegen die
|
||||
Nachrichten**anzahl** statt die höchste UID auf — mit
|
||||
`maxOpenEndedUID`-Begrenzung (statt eines naiven 2³²-1-Sentinels, der
|
||||
eine milliardenfache Schleife ausgelöst hätte) korrigiert.
|
||||
- `mail/internal/imapimport/state.go` — `Store` (Postgres,
|
||||
`mail_import_state`): persistiert `last_uidvalidity`,
|
||||
`last_synced_uid`, `interval_seconds` je Mandant/Postfach
|
||||
(Akzeptanzkriterium 3, übersteht Neustarts, da nie im
|
||||
Prozessspeicher).
|
||||
- `mail/internal/imapimport/scheduler.go` — `Scheduler.RunOnce`:
|
||||
UID-Vergleich klassifiziert Nachrichten als neu vs. bestehend
|
||||
(Akzeptanzkriterium 1), Fortschritt wird NACH JEDER einzelnen neuen
|
||||
Nachricht persistiert (nicht erst am Ende), UIDVALIDITY-Änderung löst
|
||||
vollständigen Resync aus (Akzeptanzkriterium 2, bekannten
|
||||
archivmail-Fehler UIDVALIDITY=0 vermieden).
|
||||
- `mail/internal/imapimport/client_real.go` — `RealClient`: echtes
|
||||
IMAP4rev1 über TCP (LOGIN/SELECT/UID FETCH/LOGOUT), für den
|
||||
realistischen Testpostfach-Nachweis UND als produktive Anbindung an
|
||||
jeden RFC-3501-konformen Server nutzbar.
|
||||
- Kein Umbau: `mail/internal/folderstate` (ING-05) unverändert — die
|
||||
UIDVALIDITY-Erzeugung bei echtem Ordner-Neuaufbau bleibt dort, IMP-01
|
||||
reagiert nur auf eine geänderte UIDVALIDITY, erzeugt selbst keine.
|
||||
|
||||
## Prüfungen
|
||||
|
||||
| # | Prüfung | Ergebnis |
|
||||
|---|---|---|
|
||||
| 1 | Test: zwei aufeinanderfolgende Läufe importieren keine Nachricht doppelt | **bestanden** – `TestRunOnce_TwoConsecutiveRunsNoDuplicateImport`: 3 Nachrichten im ersten Lauf real importiert, zweiter Lauf gegen unverändertes Postfach liefert real 0 neue, 3 bestehende |
|
||||
| 2 | Test: simulierter Dienst-Neustart mitten im Abgleich führt zu konsistentem Endzustand | **bestanden** – `TestRunOnce_SimulatedRestartMidSyncConsistentEndState`: Handler schlägt real nach 2 von 5 Nachrichten fehl, neuer Scheduler auf demselben persistenten Store verarbeitet real GENAU die verbleibenden 3, keine der ersten 2 erneut, `last_synced_uid` real konsistent bei 5 |
|
||||
| 3 | Test gegen Testpostfach mit realistischem Nachrichtenaufkommen | **bestanden** – `TestRunOnce_AgainstRealTestMailboxWithRealisticVolume`: echter End-zu-Ende-IMAP4rev1-Lauf (`RealClient` gegen echten laufenden ING-01-Server) mit 30 Nachrichten — alle 30 real importiert, zweiter Lauf real 0 neue/30 bestehende |
|
||||
|
||||
Zusätzlich (Akzeptanzkriterium 3, Intervallkonfiguration):
|
||||
`TestSetInterval_ConfigurableAndSurvivesRestart` — konfiguriertes
|
||||
Intervall bleibt nach simuliertem Neustart (neue Store-Instanz auf
|
||||
demselben Postgres-Zustand) real erhalten.
|
||||
|
||||
## Build/Test-Ergebnis (192.168.1.131)
|
||||
|
||||
```
|
||||
go build ./... -> clean
|
||||
go vet ./... -> clean
|
||||
golangci-lint run ./... -> 0 issues
|
||||
TEST_TENANT_DSN=... go test ./internal/imapimport/... -v -timeout 60s -> 4/4 bestanden
|
||||
TEST_TENANT_DSN=... TEST_MANTICORE_URL=... go test ./... -p 1
|
||||
-> alle 15 Pakete bestanden, keine Regression (inkl. ING-01: 6/6 weiterhin grün
|
||||
nach UID-FETCH-Erweiterung)
|
||||
```
|
||||
|
||||
## Gesamtergebnis
|
||||
|
||||
**Bestanden.** Alle drei Akzeptanzkriterien und alle drei Pflichtprüfungen
|
||||
real erfüllt. Entsperrt IMP-02, IMP-03, IMP-04, IMP-05, IMP-07, IMP-08,
|
||||
IMP-09, INT-05, UX-01.
|
||||
@@ -0,0 +1,57 @@
|
||||
# ING-05 – Prüfprotokoll: Folder-State & UIDVALIDITY-Handling
|
||||
|
||||
Voraussetzung ING-01 (Fertig). ING-05 ist die direkte Vorbedingung für
|
||||
IMP-01 (gemeinsam mit ING-01, bereits Fertig) — ohne ING-05 bleibt IMP-01
|
||||
weiterhin blockiert.
|
||||
|
||||
## Umsetzung
|
||||
|
||||
- `mail/internal/folderstate/store.go` — `Store` (Postgres,
|
||||
`mail_folder_state` + `mail_folder_state_events`, gleiches Muster wie
|
||||
`dedup`/`indexworker`/`savedsearch`):
|
||||
- `GetOrCreate`/`CurrentState`: konsistente Sicht bei parallelem Zugriff
|
||||
(Akzeptanzkriterium 2) — `INSERT ... ON CONFLICT DO NOTHING` +
|
||||
Rücklese, kein Lese-dann-Schreib-Fenster.
|
||||
- `NextUID`: vergibt UIDs atomar über `UPDATE ... RETURNING` unter
|
||||
Postgres-Zeilensperre (Akzeptanzkriterium 1/3), protokolliert jede
|
||||
Vergabe als Ereignis in derselben Transaktion.
|
||||
- `Rebuild`: simulierter Ordner-Neuaufbau — `GREATEST(uidvalidity + 1,
|
||||
jetzt_in_ns)` garantiert eine STRENG neue UIDVALIDITY, auch wenn zwei
|
||||
Neuaufbauten innerhalb derselben Nanosekunde laufen; UIDNEXT wird auf
|
||||
1 zurückgesetzt.
|
||||
- `RecordDeletion`/`Events`: Löschungen ändern UIDNEXT nicht (RFC 3501:
|
||||
UIDs werden nie wiederverwendet), alle Zustandsänderungen bleiben
|
||||
nachvollziehbar (Akzeptanzkriterium 3).
|
||||
- Bekannten Fehler vermieden (archivmail: UIDVALIDITY=0 bricht Resync bei
|
||||
nicht-konformen Servern): `newUIDValidity` erzeugt den Wert selbst
|
||||
(Unix-Nanosekunden, garantiert > 0), statt einen extern gelieferten
|
||||
Wert unbesehen zu übernehmen.
|
||||
- Kein Umbau: `mail/internal/imap` (ING-01) unverändert — `folderstate`
|
||||
ist ein eigenständiges Paket, das ING-01 künftig (IMP-01) als
|
||||
`MailboxStore`-Implementierung nutzen kann, ohne dass ING-01 selbst
|
||||
angefasst werden musste.
|
||||
|
||||
## Prüfungen
|
||||
|
||||
| # | Prüfung | Ergebnis |
|
||||
|---|---|---|
|
||||
| 1 | Automatisierter Test für UIDVALIDITY-Änderung bei simuliertem Ordner-Neuaufbau | **bestanden** – `TestRebuild_ChangesUIDValidityOnSimulatedFolderRebuild`: Ordner angelegt, UID vergeben, `Rebuild` aufgerufen — UIDVALIDITY real geändert, UIDNEXT real auf 1 zurückgesetzt, `rebuilt`-Ereignis real protokolliert |
|
||||
| 2 | Nebenläufigkeitstest: zwei Sessions auf demselben Ordner ohne Inkonsistenz | **bestanden** – `TestNextUID_ConcurrentSessionsOnSameFolderNoInconsistency`: 20 reale gleichzeitige `NextUID`-Aufrufe auf demselben Ordner, alle 20 UIDs real eindeutig, keine Dopplung |
|
||||
| 3 | Test für UIDNEXT-Monotonie über viele Einfüge-/Löschzyklen | **bestanden** – `TestNextUID_MonotonicAcrossManyInsertDeleteCycles`: 200 Zyklen, jede zweite Nachricht real "gelöscht" — UIDNEXT bleibt real strikt monoton steigend, Löschungen beeinflussen die Vergabe nicht |
|
||||
|
||||
## Build/Test-Ergebnis (192.168.1.131)
|
||||
|
||||
```
|
||||
go build ./... -> clean
|
||||
go vet ./... -> clean
|
||||
golangci-lint run ./... -> 0 issues
|
||||
TEST_TENANT_DSN=... go test ./internal/folderstate/... -v -> 3/3 bestanden
|
||||
TEST_TENANT_DSN=... TEST_MANTICORE_URL=... go test ./... -p 1
|
||||
-> alle 13 Pakete bestanden, keine Regression
|
||||
```
|
||||
|
||||
## Gesamtergebnis
|
||||
|
||||
**Bestanden.** Alle drei Akzeptanzkriterien und alle drei Pflichtprüfungen
|
||||
real erfüllt. Entsperrt IMP-01 (gemeinsam mit ING-01, bereits Fertig) und
|
||||
ING-10.
|
||||
@@ -0,0 +1,9 @@
|
||||
CREATE TABLE IF NOT EXISTS mail_folder_state (
|
||||
tenant_slug TEXT NOT NULL,
|
||||
mailbox_name TEXT NOT NULL,
|
||||
uidvalidity BIGINT NOT NULL,
|
||||
uidnext BIGINT NOT NULL,
|
||||
created_at TIMESTAMPTZ NOT NULL DEFAULT now(),
|
||||
updated_at TIMESTAMPTZ NOT NULL DEFAULT now(),
|
||||
PRIMARY KEY (tenant_slug, mailbox_name)
|
||||
)
|
||||
@@ -0,0 +1,8 @@
|
||||
CREATE TABLE IF NOT EXISTS mail_folder_state_events (
|
||||
id BIGSERIAL PRIMARY KEY,
|
||||
tenant_slug TEXT NOT NULL,
|
||||
mailbox_name TEXT NOT NULL,
|
||||
event_type TEXT NOT NULL CHECK (event_type IN ('uid_assigned', 'deleted', 'rebuilt')),
|
||||
uid BIGINT,
|
||||
created_at TIMESTAMPTZ NOT NULL DEFAULT now()
|
||||
)
|
||||
@@ -0,0 +1,251 @@
|
||||
// Package folderstate implementiert ING-05: Folder-State-Verwaltung
|
||||
// inklusive UIDVALIDITY/UIDNEXT-Handling für IMAP-Ordner (RFC 3501
|
||||
// §2.3.1.1), damit Clients (mail/internal/imap, ING-01) und
|
||||
// Importvorgänge (IMP-01) konsistente Sichten erhalten. Persistiert in
|
||||
// Postgres, gleiches Muster wie mail/internal/dedup/indexworker/
|
||||
// savedsearch — kein zentraler Migrationsläufer für Mandanten-
|
||||
// Datenbanken im Mail-Modul vorhanden, EnsureSchema legt die Tabellen
|
||||
// idempotent an.
|
||||
//
|
||||
// Bekannten Fehler vermeiden (siehe ING-01/repos-analyse-mail-reuse.md):
|
||||
// archivmail brach den Resync bei UIDVALIDITY=0 nicht-konformer Server —
|
||||
// dieses Paket erzeugt UIDVALIDITY selbst (Unix-Zeitstempel beim
|
||||
// Ordner-Neuaufbau, garantiert > 0 und monoton wachsend über
|
||||
// aufeinanderfolgende Neuaufbauten hinweg) statt einen von außen
|
||||
// gelieferten Wert unbesehen zu übernehmen.
|
||||
package folderstate
|
||||
|
||||
import (
|
||||
"context"
|
||||
_ "embed"
|
||||
"errors"
|
||||
"fmt"
|
||||
"time"
|
||||
|
||||
"github.com/jackc/pgx/v5"
|
||||
"github.com/jackc/pgx/v5/pgxpool"
|
||||
)
|
||||
|
||||
//go:embed migrations/0001_mail_folder_state.sql
|
||||
var schemaMigration string
|
||||
|
||||
//go:embed migrations/0002_mail_folder_state_events.sql
|
||||
var eventsSchemaMigration string
|
||||
|
||||
// EventType (Akzeptanzkriterium 3: State-Änderungen nachvollziehbar
|
||||
// persistiert).
|
||||
const (
|
||||
EventUIDAssigned = "uid_assigned"
|
||||
EventDeleted = "deleted"
|
||||
EventRebuilt = "rebuilt"
|
||||
)
|
||||
|
||||
// FolderState ist der aktuelle UIDVALIDITY/UIDNEXT-Zustand eines Ordners.
|
||||
type FolderState struct {
|
||||
TenantSlug string
|
||||
MailboxName string
|
||||
UIDValidity uint64
|
||||
UIDNext uint64
|
||||
}
|
||||
|
||||
// Event ist ein einzelner, nachvollziehbarer Zustandsänderungseintrag.
|
||||
type Event struct {
|
||||
EventType string
|
||||
UID *uint64
|
||||
CreatedAt time.Time
|
||||
}
|
||||
|
||||
// Store verwaltet Folder-State je Mandant und Postfach.
|
||||
type Store struct {
|
||||
pool *pgxpool.Pool
|
||||
// now ist austauschbar für Tests (deterministische UIDVALIDITY-Werte).
|
||||
now func() time.Time
|
||||
}
|
||||
|
||||
func NewStore(pool *pgxpool.Pool) *Store {
|
||||
return &Store{pool: pool, now: time.Now}
|
||||
}
|
||||
|
||||
// EnsureSchema legt die Tabellen an, falls sie noch nicht existieren.
|
||||
func (s *Store) EnsureSchema(ctx context.Context) error {
|
||||
if _, err := s.pool.Exec(ctx, schemaMigration); err != nil {
|
||||
return fmt.Errorf("folderstate: schema anlegen: %w", err)
|
||||
}
|
||||
if _, err := s.pool.Exec(ctx, eventsSchemaMigration); err != nil {
|
||||
return fmt.Errorf("folderstate: ereignis-schema anlegen: %w", err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// GetOrCreate liefert den aktuellen Zustand eines Ordners und legt ihn
|
||||
// bei erstem Zugriff neu an (UIDNEXT beginnt bei 1, RFC 3501 §2.3.1.1).
|
||||
// Konsistent bei parallelem Zugriff (Akzeptanzkriterium 2): INSERT ...
|
||||
// ON CONFLICT DO NOTHING + Rücklese, kein Lese-dann-Schreib-Fenster.
|
||||
func (s *Store) GetOrCreate(ctx context.Context, tenantSlug, mailboxName string) (FolderState, error) {
|
||||
uidvalidity := s.newUIDValidity()
|
||||
if _, err := s.pool.Exec(ctx, `
|
||||
INSERT INTO mail_folder_state (tenant_slug, mailbox_name, uidvalidity, uidnext)
|
||||
VALUES ($1, $2, $3, 1)
|
||||
ON CONFLICT (tenant_slug, mailbox_name) DO NOTHING
|
||||
`, tenantSlug, mailboxName, uidvalidity); err != nil {
|
||||
return FolderState{}, fmt.Errorf("folderstate: ordner anlegen: %w", err)
|
||||
}
|
||||
return s.CurrentState(ctx, tenantSlug, mailboxName)
|
||||
}
|
||||
|
||||
// CurrentState liest den Zustand ohne ihn anzulegen (Akzeptanzkriterium
|
||||
// 2: konsistente Sicht bei SELECT/EXAMINE).
|
||||
func (s *Store) CurrentState(ctx context.Context, tenantSlug, mailboxName string) (FolderState, error) {
|
||||
var st FolderState
|
||||
st.TenantSlug = tenantSlug
|
||||
st.MailboxName = mailboxName
|
||||
err := s.pool.QueryRow(ctx, `
|
||||
SELECT uidvalidity, uidnext FROM mail_folder_state
|
||||
WHERE tenant_slug = $1 AND mailbox_name = $2
|
||||
`, tenantSlug, mailboxName).Scan(&st.UIDValidity, &st.UIDNext)
|
||||
if err != nil {
|
||||
if errors.Is(err, pgx.ErrNoRows) {
|
||||
return FolderState{}, ErrNotFound
|
||||
}
|
||||
return FolderState{}, fmt.Errorf("folderstate: zustand lesen: %w", err)
|
||||
}
|
||||
return st, nil
|
||||
}
|
||||
|
||||
// ErrNotFound wird geliefert, wenn für den angefragten Ordner noch kein
|
||||
// Zustand existiert (GetOrCreate anlegen lassen, statt hier zu raten).
|
||||
var ErrNotFound = errors.New("folderstate: ordner nicht gefunden")
|
||||
|
||||
// NextUID vergibt atomar die nächste UID für eine neu eintreffende
|
||||
// Nachricht (Akzeptanzkriterium 1/3) und protokolliert die Vergabe.
|
||||
// Nebenläufigkeitssicher: UPDATE ... RETURNING läuft unter Postgres'
|
||||
// Zeilensperre, zwei gleichzeitige Aufrufe für denselben Ordner können
|
||||
// niemals dieselbe UID liefern (Pflichtprüfung 2).
|
||||
func (s *Store) NextUID(ctx context.Context, tenantSlug, mailboxName string) (uid uint64, err error) {
|
||||
tx, err := s.pool.Begin(ctx)
|
||||
if err != nil {
|
||||
return 0, fmt.Errorf("folderstate: transaktion starten: %w", err)
|
||||
}
|
||||
defer func() { _ = tx.Rollback(ctx) }()
|
||||
|
||||
err = tx.QueryRow(ctx, `
|
||||
UPDATE mail_folder_state
|
||||
SET uidnext = uidnext + 1, updated_at = now()
|
||||
WHERE tenant_slug = $1 AND mailbox_name = $2
|
||||
RETURNING uidnext - 1
|
||||
`, tenantSlug, mailboxName).Scan(&uid)
|
||||
if err != nil {
|
||||
if errors.Is(err, pgx.ErrNoRows) {
|
||||
return 0, ErrNotFound
|
||||
}
|
||||
return 0, fmt.Errorf("folderstate: uid vergeben: %w", err)
|
||||
}
|
||||
|
||||
if _, err := tx.Exec(ctx, `
|
||||
INSERT INTO mail_folder_state_events (tenant_slug, mailbox_name, event_type, uid)
|
||||
VALUES ($1, $2, $3, $4)
|
||||
`, tenantSlug, mailboxName, EventUIDAssigned, uid); err != nil {
|
||||
return 0, fmt.Errorf("folderstate: ereignis protokollieren: %w", err)
|
||||
}
|
||||
|
||||
if err := tx.Commit(ctx); err != nil {
|
||||
return 0, fmt.Errorf("folderstate: uid-vergabe committen: %w", err)
|
||||
}
|
||||
return uid, nil
|
||||
}
|
||||
|
||||
// RecordDeletion protokolliert die Löschung einer Nachricht mit
|
||||
// gegebener UID (Akzeptanzkriterium 3). UIDNEXT bleibt unverändert —
|
||||
// gelöschte UIDs werden gemäß RFC 3501 niemals wiederverwendet.
|
||||
func (s *Store) RecordDeletion(ctx context.Context, tenantSlug, mailboxName string, uid uint64) error {
|
||||
if _, err := s.pool.Exec(ctx, `
|
||||
INSERT INTO mail_folder_state_events (tenant_slug, mailbox_name, event_type, uid)
|
||||
VALUES ($1, $2, $3, $4)
|
||||
`, tenantSlug, mailboxName, EventDeleted, uid); err != nil {
|
||||
return fmt.Errorf("folderstate: löschung protokollieren: %w", err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// Rebuild simuliert einen Ordner-Neuaufbau (z. B. nach erkannter
|
||||
// Inkonsistenz oder bei einem Server, der seinerseits eine neue
|
||||
// UIDVALIDITY meldet): vergibt eine garantiert neue UIDVALIDITY und
|
||||
// setzt UIDNEXT zurück auf 1 (Pflichtprüfung 1).
|
||||
func (s *Store) Rebuild(ctx context.Context, tenantSlug, mailboxName string) (FolderState, error) {
|
||||
candidateUIDValidity := s.newUIDValidity()
|
||||
|
||||
tx, err := s.pool.Begin(ctx)
|
||||
if err != nil {
|
||||
return FolderState{}, fmt.Errorf("folderstate: transaktion starten: %w", err)
|
||||
}
|
||||
defer func() { _ = tx.Rollback(ctx) }()
|
||||
|
||||
// GREATEST(...)+1 garantiert eine STRENG größere UIDVALIDITY als die
|
||||
// bisherige, unabhängig von der Uhrenauflösung — zwei Neuaufbauten
|
||||
// innerhalb derselben Nanosekunde dürfen niemals denselben Wert
|
||||
// liefern (Pflichtprüfung 1).
|
||||
var newUIDValidity uint64
|
||||
err = tx.QueryRow(ctx, `
|
||||
UPDATE mail_folder_state
|
||||
SET uidvalidity = GREATEST(uidvalidity + 1, $3), uidnext = 1, updated_at = now()
|
||||
WHERE tenant_slug = $1 AND mailbox_name = $2
|
||||
RETURNING uidvalidity
|
||||
`, tenantSlug, mailboxName, candidateUIDValidity).Scan(&newUIDValidity)
|
||||
if err != nil {
|
||||
if errors.Is(err, pgx.ErrNoRows) {
|
||||
return FolderState{}, ErrNotFound
|
||||
}
|
||||
return FolderState{}, fmt.Errorf("folderstate: neuaufbau: %w", err)
|
||||
}
|
||||
|
||||
if _, err := tx.Exec(ctx, `
|
||||
INSERT INTO mail_folder_state_events (tenant_slug, mailbox_name, event_type)
|
||||
VALUES ($1, $2, $3)
|
||||
`, tenantSlug, mailboxName, EventRebuilt); err != nil {
|
||||
return FolderState{}, fmt.Errorf("folderstate: neuaufbau-ereignis protokollieren: %w", err)
|
||||
}
|
||||
|
||||
if err := tx.Commit(ctx); err != nil {
|
||||
return FolderState{}, fmt.Errorf("folderstate: neuaufbau committen: %w", err)
|
||||
}
|
||||
return FolderState{TenantSlug: tenantSlug, MailboxName: mailboxName, UIDValidity: newUIDValidity, UIDNext: 1}, nil
|
||||
}
|
||||
|
||||
// Events liefert die protokollierten Zustandsänderungen eines Ordners in
|
||||
// zeitlicher Reihenfolge (Akzeptanzkriterium 3: nachvollziehbar).
|
||||
func (s *Store) Events(ctx context.Context, tenantSlug, mailboxName string) ([]Event, error) {
|
||||
rows, err := s.pool.Query(ctx, `
|
||||
SELECT event_type, uid, created_at FROM mail_folder_state_events
|
||||
WHERE tenant_slug = $1 AND mailbox_name = $2
|
||||
ORDER BY id ASC
|
||||
`, tenantSlug, mailboxName)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("folderstate: ereignisse lesen: %w", err)
|
||||
}
|
||||
defer rows.Close()
|
||||
|
||||
var events []Event
|
||||
for rows.Next() {
|
||||
var e Event
|
||||
var uid *int64
|
||||
if err := rows.Scan(&e.EventType, &uid, &e.CreatedAt); err != nil {
|
||||
return nil, fmt.Errorf("folderstate: ereigniszeile lesen: %w", err)
|
||||
}
|
||||
if uid != nil {
|
||||
u := uint64(*uid)
|
||||
e.UID = &u
|
||||
}
|
||||
events = append(events, e)
|
||||
}
|
||||
if err := rows.Err(); err != nil {
|
||||
return nil, fmt.Errorf("folderstate: ereignisse iterieren: %w", err)
|
||||
}
|
||||
return events, nil
|
||||
}
|
||||
|
||||
// newUIDValidity erzeugt eine garantiert positive, für praktische Zwecke
|
||||
// eindeutige UIDVALIDITY (Unix-Nanosekunden) — vermeidet den bekannten
|
||||
// archivmail-Fehler UIDVALIDITY=0.
|
||||
func (s *Store) newUIDValidity() uint64 {
|
||||
return uint64(s.now().UnixNano())
|
||||
}
|
||||
@@ -0,0 +1,176 @@
|
||||
// Integrationstest (ING-05): echte Postgres-Instanz, folgt derselben
|
||||
// Testhost-Konvention wie mail/internal/dedup/indexworker/savedsearch —
|
||||
// TEST_TENANT_DSN.
|
||||
package folderstate
|
||||
|
||||
import (
|
||||
"context"
|
||||
"os"
|
||||
"sync"
|
||||
"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() {
|
||||
ctx := context.Background()
|
||||
_, _ = pool.Exec(ctx, `DELETE FROM mail_folder_state WHERE tenant_slug LIKE 'mandant-ing05-%'`)
|
||||
_, _ = pool.Exec(ctx, `DELETE FROM mail_folder_state_events WHERE tenant_slug LIKE 'mandant-ing05-%'`)
|
||||
})
|
||||
return store
|
||||
}
|
||||
|
||||
// TestRebuild_ChangesUIDValidityOnSimulatedFolderRebuild ist die
|
||||
// geforderte Pflichtprüfung 1: automatisierter Test für
|
||||
// UIDVALIDITY-Änderung bei simuliertem Ordner-Neuaufbau.
|
||||
func TestRebuild_ChangesUIDValidityOnSimulatedFolderRebuild(t *testing.T) {
|
||||
store := setupStore(t)
|
||||
ctx := context.Background()
|
||||
tenant := "mandant-ing05-rebuild"
|
||||
|
||||
initial, err := store.GetOrCreate(ctx, tenant, "INBOX")
|
||||
if err != nil {
|
||||
t.Fatalf("getorcreate: %v", err)
|
||||
}
|
||||
if initial.UIDValidity == 0 {
|
||||
t.Fatal("erwartete uidvalidity != 0 (bekannter archivmail-fehler vermeiden)")
|
||||
}
|
||||
|
||||
// UIDNEXT vor dem Neuaufbau real erhöhen, damit der Reset auf 1
|
||||
// nachweisbar ist.
|
||||
if _, err := store.NextUID(ctx, tenant, "INBOX"); err != nil {
|
||||
t.Fatalf("nextuid: %v", err)
|
||||
}
|
||||
|
||||
rebuilt, err := store.Rebuild(ctx, tenant, "INBOX")
|
||||
if err != nil {
|
||||
t.Fatalf("rebuild: %v", err)
|
||||
}
|
||||
if rebuilt.UIDValidity == initial.UIDValidity {
|
||||
t.Fatalf("erwartete geänderte uidvalidity nach neuaufbau, habe weiterhin %d", rebuilt.UIDValidity)
|
||||
}
|
||||
if rebuilt.UIDNext != 1 {
|
||||
t.Fatalf("erwartete uidnext=1 nach neuaufbau, habe %d", rebuilt.UIDNext)
|
||||
}
|
||||
|
||||
events, err := store.Events(ctx, tenant, "INBOX")
|
||||
if err != nil {
|
||||
t.Fatalf("events: %v", err)
|
||||
}
|
||||
found := false
|
||||
for _, e := range events {
|
||||
if e.EventType == EventRebuilt {
|
||||
found = true
|
||||
}
|
||||
}
|
||||
if !found {
|
||||
t.Fatal("erwartete protokolliertes 'rebuilt'-ereignis (akzeptanzkriterium 3: nachvollziehbar)")
|
||||
}
|
||||
}
|
||||
|
||||
// TestNextUID_ConcurrentSessionsOnSameFolderNoInconsistency ist die
|
||||
// geforderte Pflichtprüfung 2: Nebenläufigkeitstest — zwei Sessions auf
|
||||
// demselben Ordner ohne Inkonsistenz.
|
||||
func TestNextUID_ConcurrentSessionsOnSameFolderNoInconsistency(t *testing.T) {
|
||||
store := setupStore(t)
|
||||
ctx := context.Background()
|
||||
tenant := "mandant-ing05-concurrent"
|
||||
|
||||
if _, err := store.GetOrCreate(ctx, tenant, "INBOX"); err != nil {
|
||||
t.Fatalf("getorcreate: %v", err)
|
||||
}
|
||||
|
||||
const parallelSessions = 20
|
||||
var wg sync.WaitGroup
|
||||
uids := make(chan uint64, parallelSessions)
|
||||
errs := make(chan error, parallelSessions)
|
||||
|
||||
for i := 0; i < parallelSessions; i++ {
|
||||
wg.Add(1)
|
||||
go func() {
|
||||
defer wg.Done()
|
||||
uid, err := store.NextUID(ctx, tenant, "INBOX")
|
||||
if err != nil {
|
||||
errs <- err
|
||||
return
|
||||
}
|
||||
uids <- uid
|
||||
}()
|
||||
}
|
||||
wg.Wait()
|
||||
close(uids)
|
||||
close(errs)
|
||||
|
||||
for err := range errs {
|
||||
t.Fatalf("nextuid unter nebenläufigkeit: %v", err)
|
||||
}
|
||||
|
||||
seen := make(map[uint64]bool, parallelSessions)
|
||||
for uid := range uids {
|
||||
if seen[uid] {
|
||||
t.Fatalf("uid %d doppelt vergeben — inkonsistenz unter nebenläufigem zugriff", uid)
|
||||
}
|
||||
seen[uid] = true
|
||||
}
|
||||
if len(seen) != parallelSessions {
|
||||
t.Fatalf("erwartete %d eindeutige uids, habe %d", parallelSessions, len(seen))
|
||||
}
|
||||
}
|
||||
|
||||
// TestNextUID_MonotonicAcrossManyInsertDeleteCycles ist die geforderte
|
||||
// Pflichtprüfung 3: Test für UIDNEXT-Monotonie über viele Einfüge-/
|
||||
// Löschzyklen.
|
||||
func TestNextUID_MonotonicAcrossManyInsertDeleteCycles(t *testing.T) {
|
||||
store := setupStore(t)
|
||||
ctx := context.Background()
|
||||
tenant := "mandant-ing05-monotonie"
|
||||
|
||||
if _, err := store.GetOrCreate(ctx, tenant, "INBOX"); err != nil {
|
||||
t.Fatalf("getorcreate: %v", err)
|
||||
}
|
||||
|
||||
var lastUID uint64
|
||||
for i := 0; i < 200; i++ {
|
||||
uid, err := store.NextUID(ctx, tenant, "INBOX")
|
||||
if err != nil {
|
||||
t.Fatalf("nextuid (zyklus %d): %v", i, err)
|
||||
}
|
||||
if i > 0 && uid <= lastUID {
|
||||
t.Fatalf("uidnext nicht monoton steigend: zyklus %d, vorherige uid=%d, neue uid=%d", i, lastUID, uid)
|
||||
}
|
||||
lastUID = uid
|
||||
|
||||
// Löschung darf UIDNEXT NICHT verändern (RFC 3501: UIDs werden nie
|
||||
// wiederverwendet) — jede zweite Nachricht wird "gelöscht".
|
||||
if i%2 == 0 {
|
||||
if err := store.RecordDeletion(ctx, tenant, "INBOX", uid); err != nil {
|
||||
t.Fatalf("recorddeletion (zyklus %d): %v", i, err)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
final, err := store.CurrentState(ctx, tenant, "INBOX")
|
||||
if err != nil {
|
||||
t.Fatalf("currentstate: %v", err)
|
||||
}
|
||||
if final.UIDNext != lastUID+1 {
|
||||
t.Fatalf("erwartete uidnext=%d nach 200 vergebenen uids, habe %d", lastUID+1, final.UIDNext)
|
||||
}
|
||||
}
|
||||
@@ -53,7 +53,7 @@ func (s *Session) handleSelect(ctx context.Context, cmd command) bool {
|
||||
}
|
||||
|
||||
mailboxName := cmd.Args[0]
|
||||
exists, ok, err := s.store.Select(ctx, mailboxName)
|
||||
exists, uidvalidity, ok, err := s.store.Select(ctx, mailboxName)
|
||||
if err != nil || !ok {
|
||||
// Fehlgeschlagenes SELECT lässt den Zustand laut RFC 3501 §6.3.1
|
||||
// auf Authenticated zurückfallen, nie in Selected mit ungültigem
|
||||
@@ -65,6 +65,12 @@ func (s *Session) handleSelect(ctx context.Context, cmd command) bool {
|
||||
if err := writeUntagged(s.writer, fmt.Sprintf("%d EXISTS", exists)); err != nil {
|
||||
return false
|
||||
}
|
||||
// RFC 3501 §2.3.1.1: UIDVALIDITY ist Pflichtbestandteil der
|
||||
// SELECT-Antwort — Grundlage für IMP-01s Erkennung eines
|
||||
// Ordner-Neuaufbaus.
|
||||
if err := writeUntagged(s.writer, fmt.Sprintf("OK [UIDVALIDITY %d] UIDs valid", uidvalidity)); err != nil {
|
||||
return false
|
||||
}
|
||||
s.state = Selected
|
||||
s.mailbox = mailboxName
|
||||
s.mailboxSize = uint32(exists)
|
||||
@@ -90,13 +96,48 @@ func (s *Session) handleFetch(ctx context.Context, cmd command) bool {
|
||||
if err != nil {
|
||||
return s.writeErr(cmd.Tag, "NO", "FETCH failed")
|
||||
}
|
||||
return s.writeFetchResults(cmd.Tag, "FETCH", messages)
|
||||
}
|
||||
|
||||
// handleUIDFetch implementiert "UID FETCH" (RFC 3501 §6.4.8) — wie FETCH,
|
||||
// aber uid-set statt Sequenzsatz, Grundlage für IMP-01s UID-basierten
|
||||
// Delta-Sync.
|
||||
func (s *Session) handleUIDFetch(ctx context.Context, cmd command) bool {
|
||||
if s.state != Selected {
|
||||
return s.writeErr(cmd.Tag, "BAD", "UID FETCH not allowed in "+s.state.String()+" state")
|
||||
}
|
||||
if len(cmd.Args) < 2 {
|
||||
return s.writeErr(cmd.Tag, "BAD", "UID FETCH requires a uid set")
|
||||
}
|
||||
// "*" in einem UID-Satz bedeutet "höchste vorhandene UID", NICHT die
|
||||
// NachrichtenANZAHL (s.mailboxSize) — UIDs können durch Löschungen
|
||||
// weit über der Nachrichtenzahl liegen (siehe mail/internal/
|
||||
// folderstate, ING-05: UIDs werden nie wiederverwendet). Da
|
||||
// parseSequenceSet einen Bereich materialisiert, wird "*" hier auf
|
||||
// maxOpenEndedUID begrenzt statt auf 2^32-1 — verhindert eine
|
||||
// Milliarden Einträge lange Schleife bei einem einzelnen offenen
|
||||
// Bereich. FetchByUID liefert ohnehin nur tatsächlich vorhandene
|
||||
// UIDs zurück, die Begrenzung ist für reale Postfachgrößen harmlos.
|
||||
uidSet, err := parseSequenceSet(cmd.Args[1], maxOpenEndedUID)
|
||||
if err != nil {
|
||||
return s.writeErr(cmd.Tag, "BAD", "UID FETCH: invalid uid set")
|
||||
}
|
||||
|
||||
messages, err := s.store.FetchByUID(ctx, s.mailbox, uidSet)
|
||||
if err != nil {
|
||||
return s.writeErr(cmd.Tag, "NO", "UID FETCH failed")
|
||||
}
|
||||
return s.writeFetchResults(cmd.Tag, "UID FETCH", messages)
|
||||
}
|
||||
|
||||
func (s *Session) writeFetchResults(tag, completedText string, messages []Message) bool {
|
||||
for _, m := range messages {
|
||||
text := fmt.Sprintf("%d FETCH (FLAGS (%s))", m.SequenceNumber, strings.Join(m.Flags, " "))
|
||||
text := fmt.Sprintf("%d FETCH (UID %d FLAGS (%s))", m.SequenceNumber, m.UID, strings.Join(m.Flags, " "))
|
||||
if err := writeUntagged(s.writer, text); err != nil {
|
||||
return false
|
||||
}
|
||||
}
|
||||
return s.writeErr(cmd.Tag, "OK", "FETCH completed")
|
||||
return s.writeErr(tag, "OK", completedText+" completed")
|
||||
}
|
||||
|
||||
// handleLogout ist in jedem Zustand erlaubt und beendet die Sitzung.
|
||||
@@ -114,6 +155,11 @@ func (s *Session) handleLogout(cmd command) bool {
|
||||
// tatsächliche Nachrichtenzahl des gewählten Postfachs, von SELECT
|
||||
// gemeldet). Volle RFC-3501-Sequenzsatz-Grammatik (verschachtelte
|
||||
// Bereiche etc.) ist bewusst nicht Bestandteil dieser kleinsten Lösung.
|
||||
// maxOpenEndedUID begrenzt, wie weit ein offener UID-Bereich ("N:*")
|
||||
// materialisiert wird — deckt reale Postfachgrößen komfortabel ab, ohne
|
||||
// bei einem einzelnen Kommando Milliarden Slice-Einträge zu erzeugen.
|
||||
const maxOpenEndedUID = 1_000_000
|
||||
|
||||
func parseSequenceSet(raw string, maxSeq uint32) ([]uint32, error) {
|
||||
var result []uint32
|
||||
for _, part := range strings.Split(raw, ",") {
|
||||
|
||||
@@ -28,9 +28,9 @@ type fakeMailboxStore struct {
|
||||
mailboxes map[string][]Message
|
||||
}
|
||||
|
||||
func (f fakeMailboxStore) Select(_ context.Context, mailboxName string) (int, bool, error) {
|
||||
func (f fakeMailboxStore) Select(_ context.Context, mailboxName string) (int, uint64, bool, error) {
|
||||
msgs, ok := f.mailboxes[mailboxName]
|
||||
return len(msgs), ok, nil
|
||||
return len(msgs), 1, ok, nil
|
||||
}
|
||||
|
||||
func (f fakeMailboxStore) Fetch(_ context.Context, mailboxName string, seqNumbers []uint32) ([]Message, error) {
|
||||
@@ -51,13 +51,31 @@ func (f fakeMailboxStore) Fetch(_ context.Context, mailboxName string, seqNumber
|
||||
return result, nil
|
||||
}
|
||||
|
||||
func (f fakeMailboxStore) FetchByUID(_ context.Context, mailboxName string, uids []uint32) ([]Message, error) {
|
||||
msgs, ok := f.mailboxes[mailboxName]
|
||||
if !ok {
|
||||
return nil, errors.New("imap: postfach nicht gefunden")
|
||||
}
|
||||
wanted := make(map[uint32]bool, len(uids))
|
||||
for _, u := range uids {
|
||||
wanted[u] = true
|
||||
}
|
||||
var result []Message
|
||||
for _, m := range msgs {
|
||||
if wanted[m.UID] {
|
||||
result = append(result, m)
|
||||
}
|
||||
}
|
||||
return result, nil
|
||||
}
|
||||
|
||||
func startTestServer(t *testing.T) (addr string, stop func()) {
|
||||
t.Helper()
|
||||
auth := fakeAuthenticator{users: map[string]string{"alice": "geheim123"}}
|
||||
store := fakeMailboxStore{mailboxes: map[string][]Message{
|
||||
"INBOX": {
|
||||
{SequenceNumber: 1, Flags: []string{"\\Seen"}},
|
||||
{SequenceNumber: 2, Flags: []string{}},
|
||||
{SequenceNumber: 1, UID: 101, Flags: []string{"\\Seen"}},
|
||||
{SequenceNumber: 2, UID: 102, Flags: []string{}},
|
||||
},
|
||||
}}
|
||||
srv := NewServer(auth, store)
|
||||
@@ -203,7 +221,7 @@ func TestCommands_AllBaseCommandsAnswered(t *testing.T) {
|
||||
}
|
||||
|
||||
_, lines = c.sendTagged(t, "FETCH 1 (FLAGS)")
|
||||
if !containsSubstring(lines, "FETCH (FLAGS") {
|
||||
if !containsSubstring(lines, "FETCH (UID") {
|
||||
t.Fatalf("FETCH: erwartete FLAGS-Antwort, habe: %v", lines)
|
||||
}
|
||||
|
||||
@@ -303,3 +321,23 @@ func TestServer_50ParallelSessionsNoLeak(t *testing.T) {
|
||||
t.Errorf("parallele sitzung fehlgeschlagen: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
// TestCommands_UIDFetchReturnsUID belegt die für IMP-01 nötige
|
||||
// UID-FETCH-Erweiterung: reale UID-basierte Abfrage über echtes TCP.
|
||||
func TestCommands_UIDFetchReturnsUID(t *testing.T) {
|
||||
addr, stop := startTestServer(t)
|
||||
defer stop()
|
||||
c := dial(t, addr)
|
||||
defer c.close()
|
||||
|
||||
c.sendTagged(t, "LOGIN alice geheim123")
|
||||
c.sendTagged(t, "SELECT INBOX")
|
||||
|
||||
_, lines := c.sendTagged(t, "UID FETCH 101:102 (FLAGS)")
|
||||
if !containsSubstring(lines, "UID 101") || !containsSubstring(lines, "UID 102") {
|
||||
t.Fatalf("erwartete beide UIDs in der antwort, habe: %v", lines)
|
||||
}
|
||||
if !strings.Contains(lines[len(lines)-1], "OK") {
|
||||
t.Fatalf("erwartete OK-abschluss, habe: %v", lines)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -21,19 +21,28 @@ type Authenticator interface {
|
||||
Authenticate(ctx context.Context, username, password string) (ok bool, err error)
|
||||
}
|
||||
|
||||
// Message ist eine minimale Nachrichtendarstellung für FETCH (nur Flags,
|
||||
// keine Inhalte — Inhaltszugriff ist Sache späterer Kacheln).
|
||||
// Message ist eine minimale Nachrichtendarstellung für FETCH (nur UID +
|
||||
// Flags, keine Inhalte — Inhaltszugriff ist Sache späterer Kacheln). UID
|
||||
// wird seit IMP-01 zusätzlich zur Sequenznummer geführt (RFC 3501 §2.3.1,
|
||||
// UID FETCH) — Grundlage für IMP-01s UID-basierten Delta-Sync.
|
||||
type Message struct {
|
||||
SequenceNumber uint32
|
||||
UID uint32
|
||||
Flags []string
|
||||
}
|
||||
|
||||
// MailboxStore liefert Postfachzustand für SELECT/FETCH.
|
||||
type MailboxStore interface {
|
||||
// Select liefert die Anzahl der Nachrichten im Postfach mailboxName.
|
||||
// ok=false, wenn das Postfach nicht existiert.
|
||||
Select(ctx context.Context, mailboxName string) (exists int, ok bool, err error)
|
||||
// Select liefert die Anzahl der Nachrichten sowie die UIDVALIDITY
|
||||
// (RFC 3501 §2.3.1.1 — Pflichtbestandteil der SELECT-Antwort, Basis
|
||||
// für IMP-01s Erkennung eines Ordner-Neuaufbaus) des Postfachs
|
||||
// mailboxName. ok=false, wenn das Postfach nicht existiert.
|
||||
Select(ctx context.Context, mailboxName string) (exists int, uidvalidity uint64, ok bool, err error)
|
||||
// Fetch liefert die Nachrichten im aktuell gewählten Postfach, deren
|
||||
// Sequenznummer in seqNumbers enthalten ist.
|
||||
Fetch(ctx context.Context, mailboxName string, seqNumbers []uint32) ([]Message, error)
|
||||
// FetchByUID liefert die Nachrichten im aktuell gewählten Postfach,
|
||||
// deren UID in uids enthalten ist (RFC 3501 §6.4.8, UID FETCH) — Basis
|
||||
// für IMP-01s UID-Vergleich.
|
||||
FetchByUID(ctx context.Context, mailboxName string, uids []uint32) ([]Message, error)
|
||||
}
|
||||
|
||||
@@ -103,6 +103,11 @@ func (s *Session) dispatch(ctx context.Context, cmd command) bool {
|
||||
return s.handleSelect(ctx, cmd)
|
||||
case "FETCH":
|
||||
return s.handleFetch(ctx, cmd)
|
||||
case "UID":
|
||||
if len(cmd.Args) < 1 || strings.ToUpper(cmd.Args[0]) != "FETCH" {
|
||||
return s.writeErr(cmd.Tag, "BAD", "Unsupported UID subcommand")
|
||||
}
|
||||
return s.handleUIDFetch(ctx, cmd)
|
||||
case "LOGOUT":
|
||||
return s.handleLogout(cmd)
|
||||
default:
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
@@ -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)
|
||||
)
|
||||
@@ -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
|
||||
}
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
@@ -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
|
||||
}
|
||||
Reference in New Issue
Block a user