Compare commits
3
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
7238568918 | ||
|
|
03c47d98f4 | ||
|
|
e9947b1e28 |
@@ -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,48 @@
|
|||||||
|
# IMP-02 – Prüfprotokoll: Anhangsverarbeitung bei Import
|
||||||
|
|
||||||
|
Voraussetzung ING-04, IMP-01 (beide Fertig).
|
||||||
|
|
||||||
|
## Umsetzung
|
||||||
|
|
||||||
|
- `mail/internal/mimeparse/tolerant.go` — additive Erweiterung von ING-04
|
||||||
|
(Parse/parseMultipart bleiben UNVERÄNDERT): `ParseTolerant` bricht bei
|
||||||
|
einem einzelnen fehlerhaften Teil NICHT die gesamte Nachricht ab
|
||||||
|
(Akzeptanzkriterium 3), sondern verzeichnet ihn in `[]PartError` und
|
||||||
|
verarbeitet die übrigen Teile weiter. Setzt zusätzlich ein
|
||||||
|
Gesamtgrößenbudget über alle Teile durch (`ErrMessageTooLarge`,
|
||||||
|
Akzeptanzkriterium 2 — ergänzt das bereits vorhandene
|
||||||
|
Je-Anhang-Limit aus ING-04 um ein Je-Nachricht-Limit).
|
||||||
|
- `mail/internal/attachments/attachments.go` — `Extract`: liefert
|
||||||
|
`Attachment{Filename, Size, DeclaredContentType, VerifiedContentType}`
|
||||||
|
je Anhang (Akzeptanzkriterium 1) — `VerifiedContentType` kommt aus
|
||||||
|
`net/http.DetectContentType` (echtes Sniffing der Bytes), nicht aus der
|
||||||
|
ungeprüft übernommenen Absenderbehauptung. `Options{MaxAttachmentSize,
|
||||||
|
MaxMessageSize}` mit sinnvollen Vorgabewerten (25 MiB je Anhang,
|
||||||
|
100 MiB je Nachricht).
|
||||||
|
- Kein Umbau: `mail/internal/mimeparse` Parse/parseMultipart (ING-04)
|
||||||
|
unverändert — bestehende Tests laufen unangetastet weiter.
|
||||||
|
|
||||||
|
## Prüfungen
|
||||||
|
|
||||||
|
| # | Prüfung | Ergebnis |
|
||||||
|
|---|---|---|
|
||||||
|
| 1 | Test mit Nachricht, die einen überdimensionierten Anhang enthält, wird korrekt begrenzt | **bestanden** – `TestExtract_OversizedAttachmentIsCorrectlyLimited`: Anhang über dem Limit wird real übersprungen (nicht extrahiert), Nachrichtentext bleibt real unangetastet |
|
||||||
|
| 2 | Test mit mehreren Anhängen unterschiedlichen Typs importiert alle korrekt | **bestanden** – `TestExtract_MultipleAttachmentDifferentTypesAllImported`: PDF + PNG in einer Nachricht, beide real extrahiert, PNG-Anhang liefert real den korrekten gesniffeten Content-Type `image/png` (echte Magic-Bytes) |
|
||||||
|
| 3 | Test: ein defekter Anhang lässt Text und übrige Anhänge unangetastet | **bestanden** – `TestExtract_BrokenAttachmentLeavesTextAndOthersUntouched`: ungültiges Base64 in einem Anhang, Nachrichtentext UND der zweite, gültige Anhang kommen real unverändert an |
|
||||||
|
|
||||||
|
## Build/Test-Ergebnis (192.168.1.131)
|
||||||
|
|
||||||
|
```
|
||||||
|
go build ./... -> clean
|
||||||
|
go vet ./... -> clean
|
||||||
|
golangci-lint run ./... -> 0 issues
|
||||||
|
go test ./internal/attachments/... -v -> 3/3 bestanden
|
||||||
|
TEST_TENANT_DSN=... TEST_MANTICORE_URL=... go test ./... -p 1
|
||||||
|
-> alle 16 Pakete bestanden, keine Regression (mimeparse: 6/6 weiterhin grün
|
||||||
|
nach additiver ParseTolerant-Erweiterung)
|
||||||
|
```
|
||||||
|
|
||||||
|
## Gesamtergebnis
|
||||||
|
|
||||||
|
**Bestanden.** Alle drei Akzeptanzkriterien und alle drei Pflichtprüfungen
|
||||||
|
real erfüllt. Entsperrt IMP-06, trägt (gemeinsam mit IMP-03) zu IMP-09 bei.
|
||||||
@@ -0,0 +1,58 @@
|
|||||||
|
# IMP-04 – Prüfprotokoll: Fehlerbehandlung nicht-konformer Server
|
||||||
|
|
||||||
|
Voraussetzung IMP-01 (Fertig).
|
||||||
|
|
||||||
|
## Umsetzung
|
||||||
|
|
||||||
|
- `mail/internal/imapimport/client_real.go` erweitert:
|
||||||
|
- `resolveUIDValidity`: eine gemeldete `UIDVALIDITY=0` (bekannte
|
||||||
|
Abweichung nicht-konformer Server, known-issues-archivmail.md #5)
|
||||||
|
oder eine ganz fehlende UIDVALIDITY-Angabe löst KEINEN Abbruch mehr
|
||||||
|
aus, sondern einen definierten Fallback (Akzeptanzkriterium 1):
|
||||||
|
`fallbackUIDValidity` leitet deterministisch (FNV-1a, gleiche Technik
|
||||||
|
wie `search.DocumentID`) einen von 0 verschiedenen Ersatzwert aus dem
|
||||||
|
Postfachnamen ab — bei wiederholten Läufen gegen denselben
|
||||||
|
nicht-konformen Server bleibt der Fallback STABIL, kein unnötiger
|
||||||
|
Voll-Resync bei jedem einzelnen Lauf.
|
||||||
|
- `parseFetchLines`/`parseSingleFetchLine`: eine einzelne unerwartete
|
||||||
|
oder kaputte `FETCH`-Zeile wird protokolliert und übersprungen, alle
|
||||||
|
übrigen, korrekt lesbaren Nachrichten werden trotzdem geliefert
|
||||||
|
(Akzeptanzkriterium 2) — der gesamte Lauf bricht dafür nicht ab.
|
||||||
|
- `Logger`/`RealClient.WithLogger`: jede erkannte Abweichung läuft über
|
||||||
|
ein protokollierbares, austauschbares Logging-Ziel mit festem,
|
||||||
|
durchsuchbarem Präfix (Akzeptanzkriterium 3: für Support
|
||||||
|
nachvollziehbar) — Standard ist `log.Printf`.
|
||||||
|
- Dabei einen echten, durch die neue Logging-Logik selbst eingeführten
|
||||||
|
Bug gefunden und behoben: die getaggte Kommando-Abschlusszeile (z. B.
|
||||||
|
`"C3 OK UID FETCH completed"`) enthält ebenfalls die Zeichenfolge
|
||||||
|
`"FETCH "` und wurde beim ersten Anlauf fälschlich als "unerwartete
|
||||||
|
Serverantwort" geloggt — behoben, indem nur echte Untagged-Zeilen
|
||||||
|
(Präfix `"* "`) überhaupt als FETCH-Zeile in Betracht gezogen werden.
|
||||||
|
- Kein Umbau: `mail/internal/imap` (ING-01)/`folderstate` (ING-05)/
|
||||||
|
`imapimport/scheduler.go` (IMP-01) unverändert — IMP-04 erweitert
|
||||||
|
ausschließlich `client_real.go`.
|
||||||
|
|
||||||
|
## Prüfungen
|
||||||
|
|
||||||
|
| # | Prüfung | Ergebnis |
|
||||||
|
|---|---|---|
|
||||||
|
| 1 | Test simuliert Server mit UIDVALIDITY=0 und bestätigt greifenden Fallback | **bestanden** – `TestResolveUIDValidity_ZeroTriggersDefinedFallbackNotAbort`: hand-gesteuerter Fake-Server meldet real `UIDVALIDITY=0`, `Sync` schlägt real NICHT fehl, liefert real einen von 0 verschiedenen, deterministischen Fallback-Wert und alle 3 Nachrichten, Fallback-Hinweis real protokolliert |
|
||||||
|
| 2 | Test mit unerwarteter/kaputter Serverantwort bestätigt Weiterlauf für übrige Nachrichten | **bestanden** – `TestParseFetchLines_UnexpectedResponseSkippedRestContinue`: 2 bewusst kaputte Zeilen zwischen 2 korrekten real gesendet — `Sync` liefert real trotzdem beide korrekt lesbaren Nachrichten, beide kaputten Zeilen real protokolliert und übersprungen, kein Abbruch |
|
||||||
|
| 3 | Regressionstest verhindert Wiederauftreten des UIDVALIDITY-Bugs | **bestanden** – `TestResolveUIDValidity_RegressionGuardAgainstZeroAbort`: direkter, vom Netzwerkpfad unabhängiger Test von `resolveUIDValidity` mit `UIDVALIDITY=0` UND mit gänzlich fehlender Angabe — beide liefern real keinen Fehler und einen Fallback-Wert != 0 |
|
||||||
|
|
||||||
|
## 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 -> 7/7 bestanden
|
||||||
|
TEST_TENANT_DSN=... TEST_MANTICORE_URL=... go test ./... -p 1
|
||||||
|
-> alle 15 Pakete bestanden, keine Regression
|
||||||
|
```
|
||||||
|
|
||||||
|
## Gesamtergebnis
|
||||||
|
|
||||||
|
**Bestanden.** Alle drei Akzeptanzkriterien und alle drei Pflichtprüfungen
|
||||||
|
real erfüllt. Entsperrt IMP-08 (gemeinsam mit QA-02, bleibt weiterhin
|
||||||
|
blockiert bis dessen übrige Abhängigkeiten fertig sind).
|
||||||
@@ -0,0 +1,109 @@
|
|||||||
|
// Package attachments implementiert IMP-02: Anhänge aus importierten
|
||||||
|
// Nachrichten extrahieren, validieren und für die Weiterverarbeitung
|
||||||
|
// (Speicherung, Virenscan — beides spätere Kacheln, siehe "Nicht
|
||||||
|
// Bestandteil dieser Kachel") bereitstellen. Baut auf ING-04
|
||||||
|
// (mail/internal/mimeparse) auf, unverändert wiederverwendet über die
|
||||||
|
// additive Erweiterung mimeparse.ParseTolerant — kein Umbau der
|
||||||
|
// bestehenden, fertigen ING-04-Logik.
|
||||||
|
package attachments
|
||||||
|
|
||||||
|
import (
|
||||||
|
"io"
|
||||||
|
"net/http"
|
||||||
|
|
||||||
|
"gitea.perlbach24.de/scripte/nexarch/mail/internal/mimeparse"
|
||||||
|
)
|
||||||
|
|
||||||
|
// Attachment ist EIN extrahierter, validierter Anhang
|
||||||
|
// (Akzeptanzkriterium 1: Originaldateiname, Größe, geprüfter
|
||||||
|
// Content-Type).
|
||||||
|
type Attachment struct {
|
||||||
|
Filename string
|
||||||
|
Size int64
|
||||||
|
Content []byte
|
||||||
|
// DeclaredContentType kommt unverändert aus dem MIME-Header des
|
||||||
|
// Absenders — NICHT vertrauenswürdig, ein Absender kann hier
|
||||||
|
// beliebiges behaupten.
|
||||||
|
DeclaredContentType string
|
||||||
|
// VerifiedContentType wird aus den tatsächlichen Bytes gesniffed
|
||||||
|
// (net/http.DetectContentType, RFC-basierte Inhaltserkennung) —
|
||||||
|
// Akzeptanzkriterium 1: "geprüfter Content-Type", unabhängig von der
|
||||||
|
// Absenderbehauptung.
|
||||||
|
VerifiedContentType string
|
||||||
|
}
|
||||||
|
|
||||||
|
// SkippedPart beschreibt einen Anhang/Teil, der NICHT extrahiert werden
|
||||||
|
// konnte — der Rest der Nachricht (Text und übrige Anhänge) bleibt davon
|
||||||
|
// unangetastet (Akzeptanzkriterium 3).
|
||||||
|
type SkippedPart struct {
|
||||||
|
Filename string
|
||||||
|
Reason error
|
||||||
|
}
|
||||||
|
|
||||||
|
// Result ist das Ergebnis einer Anhangsextraktion.
|
||||||
|
type Result struct {
|
||||||
|
Attachments []Attachment
|
||||||
|
// TextParts sind die Nicht-Anhang-Teile (Nachrichtentext) —
|
||||||
|
// unverändert aus mimeparse übernommen, diese Kachel fasst sie nicht
|
||||||
|
// an.
|
||||||
|
TextParts []mimeparse.Part
|
||||||
|
Skipped []SkippedPart
|
||||||
|
}
|
||||||
|
|
||||||
|
// DefaultMaxAttachmentSize/DefaultMaxMessageSize sind Vorgabewerte,
|
||||||
|
// überschreibbar über Options — großzügig für typische Geschäftspost
|
||||||
|
// (kleinste Lösung, keine Konfigurationsoberfläche in dieser Kachel).
|
||||||
|
const (
|
||||||
|
DefaultMaxAttachmentSize = 25 * 1024 * 1024 // 25 MiB je Anhang
|
||||||
|
DefaultMaxMessageSize = 100 * 1024 * 1024 // 100 MiB je Nachricht gesamt
|
||||||
|
)
|
||||||
|
|
||||||
|
// Options steuert die Größenlimits (Akzeptanzkriterium 2).
|
||||||
|
type Options struct {
|
||||||
|
MaxAttachmentSize int64
|
||||||
|
MaxMessageSize int64
|
||||||
|
}
|
||||||
|
|
||||||
|
func (o Options) withDefaults() Options {
|
||||||
|
if o.MaxAttachmentSize <= 0 {
|
||||||
|
o.MaxAttachmentSize = DefaultMaxAttachmentSize
|
||||||
|
}
|
||||||
|
if o.MaxMessageSize <= 0 {
|
||||||
|
o.MaxMessageSize = DefaultMaxMessageSize
|
||||||
|
}
|
||||||
|
return o
|
||||||
|
}
|
||||||
|
|
||||||
|
// Extract zerlegt eine E-Mail (RFC 5322 + MIME) in Anhänge und
|
||||||
|
// Textteile. Ein einzelner fehlerhafter oder überdimensionierter Anhang
|
||||||
|
// blockiert NICHT die Verarbeitung der übrigen Teile
|
||||||
|
// (Akzeptanzkriterium 3) — nur eine strukturell unlesbare Nachricht
|
||||||
|
// (kaputte Kopfzeilen) liefert einen echten Fehler.
|
||||||
|
func Extract(r io.Reader, opts Options) (Result, error) {
|
||||||
|
opts = opts.withDefaults()
|
||||||
|
|
||||||
|
msg, partErrors, err := mimeparse.ParseTolerant(r, opts.MaxAttachmentSize, opts.MaxMessageSize)
|
||||||
|
if err != nil {
|
||||||
|
return Result{}, err
|
||||||
|
}
|
||||||
|
|
||||||
|
var result Result
|
||||||
|
for _, pe := range partErrors {
|
||||||
|
result.Skipped = append(result.Skipped, SkippedPart{Filename: pe.Filename, Reason: pe.Err})
|
||||||
|
}
|
||||||
|
|
||||||
|
for _, part := range msg.Parts {
|
||||||
|
if !part.IsAttachment {
|
||||||
|
result.TextParts = append(result.TextParts, part)
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
result.Attachments = append(result.Attachments, Attachment{
|
||||||
|
Filename: part.Filename,
|
||||||
|
Size: part.Size,
|
||||||
|
Content: part.Content,
|
||||||
|
DeclaredContentType: part.ContentType,
|
||||||
|
VerifiedContentType: http.DetectContentType(part.Content),
|
||||||
|
})
|
||||||
|
}
|
||||||
|
return result, nil
|
||||||
|
}
|
||||||
@@ -0,0 +1,139 @@
|
|||||||
|
package attachments
|
||||||
|
|
||||||
|
import (
|
||||||
|
"encoding/base64"
|
||||||
|
"strings"
|
||||||
|
"testing"
|
||||||
|
)
|
||||||
|
|
||||||
|
// TestExtract_OversizedAttachmentIsCorrectlyLimited ist die geforderte
|
||||||
|
// Pflichtprüfung 1: Nachricht mit überdimensioniertem Anhang wird
|
||||||
|
// korrekt begrenzt.
|
||||||
|
func TestExtract_OversizedAttachmentIsCorrectlyLimited(t *testing.T) {
|
||||||
|
oversized := strings.Repeat("A", 200)
|
||||||
|
raw := "From: a@example.com\r\n" +
|
||||||
|
"To: b@example.com\r\n" +
|
||||||
|
"Subject: Test\r\n" +
|
||||||
|
"MIME-Version: 1.0\r\n" +
|
||||||
|
"Content-Type: multipart/mixed; boundary=\"b\"\r\n\r\n" +
|
||||||
|
"--b\r\n" +
|
||||||
|
"Content-Type: text/plain; charset=utf-8\r\n\r\n" +
|
||||||
|
"Kurzer Nachrichtentext\r\n" +
|
||||||
|
"--b\r\n" +
|
||||||
|
"Content-Type: application/octet-stream\r\n" +
|
||||||
|
"Content-Disposition: attachment; filename=\"riesig.bin\"\r\n\r\n" +
|
||||||
|
oversized + "\r\n" +
|
||||||
|
"--b--\r\n"
|
||||||
|
|
||||||
|
result, err := Extract(strings.NewReader(raw), Options{MaxAttachmentSize: 50, MaxMessageSize: DefaultMaxMessageSize})
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("extract: %v", err)
|
||||||
|
}
|
||||||
|
if len(result.Attachments) != 0 {
|
||||||
|
t.Fatalf("erwartete 0 extrahierte anhänge (überdimensioniert), habe %d", len(result.Attachments))
|
||||||
|
}
|
||||||
|
if len(result.Skipped) != 1 || result.Skipped[0].Filename != "riesig.bin" {
|
||||||
|
t.Fatalf("erwartete genau 1 übersprungenen anhang 'riesig.bin', habe: %+v", result.Skipped)
|
||||||
|
}
|
||||||
|
if len(result.TextParts) != 1 || string(result.TextParts[0].Content) != "Kurzer Nachrichtentext" {
|
||||||
|
t.Fatalf("erwartete unangetasteten text trotz überdimensioniertem anhang, habe: %+v", result.TextParts)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestExtract_MultipleAttachmentDifferentTypesAllImported ist die
|
||||||
|
// geforderte Pflichtprüfung 2: mehrere Anhänge unterschiedlichen Typs
|
||||||
|
// werden alle korrekt importiert.
|
||||||
|
func TestExtract_MultipleAttachmentDifferentTypesAllImported(t *testing.T) {
|
||||||
|
pdfContent := base64.StdEncoding.EncodeToString([]byte("%PDF-1.4 fake pdf bytes"))
|
||||||
|
pngContent := base64.StdEncoding.EncodeToString([]byte{0x89, 'P', 'N', 'G', 0x0D, 0x0A, 0x1A, 0x0A, 0, 0, 0})
|
||||||
|
|
||||||
|
raw := "From: a@example.com\r\n" +
|
||||||
|
"To: b@example.com\r\n" +
|
||||||
|
"Subject: Test\r\n" +
|
||||||
|
"MIME-Version: 1.0\r\n" +
|
||||||
|
"Content-Type: multipart/mixed; boundary=\"b\"\r\n\r\n" +
|
||||||
|
"--b\r\n" +
|
||||||
|
"Content-Type: text/plain; charset=utf-8\r\n\r\n" +
|
||||||
|
"Anbei zwei Anhänge\r\n" +
|
||||||
|
"--b\r\n" +
|
||||||
|
"Content-Type: application/pdf\r\n" +
|
||||||
|
"Content-Disposition: attachment; filename=\"rechnung.pdf\"\r\n" +
|
||||||
|
"Content-Transfer-Encoding: base64\r\n\r\n" +
|
||||||
|
pdfContent + "\r\n" +
|
||||||
|
"--b\r\n" +
|
||||||
|
"Content-Type: image/png\r\n" +
|
||||||
|
"Content-Disposition: attachment; filename=\"logo.png\"\r\n" +
|
||||||
|
"Content-Transfer-Encoding: base64\r\n\r\n" +
|
||||||
|
pngContent + "\r\n" +
|
||||||
|
"--b--\r\n"
|
||||||
|
|
||||||
|
result, err := Extract(strings.NewReader(raw), Options{})
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("extract: %v", err)
|
||||||
|
}
|
||||||
|
if len(result.Attachments) != 2 {
|
||||||
|
t.Fatalf("erwartete 2 extrahierte anhänge, habe %d: %+v", len(result.Attachments), result.Attachments)
|
||||||
|
}
|
||||||
|
byName := map[string]Attachment{}
|
||||||
|
for _, a := range result.Attachments {
|
||||||
|
byName[a.Filename] = a
|
||||||
|
}
|
||||||
|
pdf, ok := byName["rechnung.pdf"]
|
||||||
|
if !ok || pdf.DeclaredContentType != "application/pdf" {
|
||||||
|
t.Fatalf("pdf-anhang fehlt oder falscher deklarierter typ: %+v", byName)
|
||||||
|
}
|
||||||
|
if !strings.Contains(pdf.VerifiedContentType, "text/plain") && !strings.Contains(pdf.VerifiedContentType, "application/") {
|
||||||
|
// http.DetectContentType erkennt unser Fake-PDF (kein echter PDF-
|
||||||
|
// Header) plausibel als Text — hier zählt nur, dass überhaupt ein
|
||||||
|
// echter, aus dem Inhalt gesniffter Wert vorliegt (Akzeptanz-
|
||||||
|
// kriterium 1: geprüfter statt blind übernommener Content-Type).
|
||||||
|
t.Fatalf("erwartete real gesniffeden content-type, habe: %q", pdf.VerifiedContentType)
|
||||||
|
}
|
||||||
|
png, ok := byName["logo.png"]
|
||||||
|
if !ok || png.DeclaredContentType != "image/png" {
|
||||||
|
t.Fatalf("png-anhang fehlt oder falscher deklarierter typ: %+v", byName)
|
||||||
|
}
|
||||||
|
if png.VerifiedContentType != "image/png" {
|
||||||
|
t.Fatalf("erwartete real gesniffeten content-type image/png (echte PNG-Magic-Bytes), habe: %q", png.VerifiedContentType)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestExtract_BrokenAttachmentLeavesTextAndOthersUntouched ist die
|
||||||
|
// geforderte Pflichtprüfung 3: ein defekter Anhang lässt Text und übrige
|
||||||
|
// Anhänge unangetastet.
|
||||||
|
func TestExtract_BrokenAttachmentLeavesTextAndOthersUntouched(t *testing.T) {
|
||||||
|
goodContent := base64.StdEncoding.EncodeToString([]byte("echter anhangsinhalt"))
|
||||||
|
raw := "From: a@example.com\r\n" +
|
||||||
|
"To: b@example.com\r\n" +
|
||||||
|
"Subject: Test\r\n" +
|
||||||
|
"MIME-Version: 1.0\r\n" +
|
||||||
|
"Content-Type: multipart/mixed; boundary=\"b\"\r\n\r\n" +
|
||||||
|
"--b\r\n" +
|
||||||
|
"Content-Type: text/plain; charset=utf-8\r\n\r\n" +
|
||||||
|
"Wichtiger Nachrichtentext\r\n" +
|
||||||
|
"--b\r\n" +
|
||||||
|
"Content-Type: application/octet-stream\r\n" +
|
||||||
|
"Content-Disposition: attachment; filename=\"kaputt.bin\"\r\n" +
|
||||||
|
"Content-Transfer-Encoding: base64\r\n\r\n" +
|
||||||
|
"DAS_IST_KEIN_GUELTIGES_BASE64!!!\r\n" +
|
||||||
|
"--b\r\n" +
|
||||||
|
"Content-Type: application/octet-stream\r\n" +
|
||||||
|
"Content-Disposition: attachment; filename=\"gut.bin\"\r\n" +
|
||||||
|
"Content-Transfer-Encoding: base64\r\n\r\n" +
|
||||||
|
goodContent + "\r\n" +
|
||||||
|
"--b--\r\n"
|
||||||
|
|
||||||
|
result, err := Extract(strings.NewReader(raw), Options{})
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("extract: %v", err)
|
||||||
|
}
|
||||||
|
if len(result.TextParts) != 1 || string(result.TextParts[0].Content) != "Wichtiger Nachrichtentext" {
|
||||||
|
t.Fatalf("erwartete unangetasteten text trotz defektem anhang, habe: %+v", result.TextParts)
|
||||||
|
}
|
||||||
|
if len(result.Attachments) != 1 || result.Attachments[0].Filename != "gut.bin" {
|
||||||
|
t.Fatalf("erwartete den guten anhang unangetastet, habe: %+v", result.Attachments)
|
||||||
|
}
|
||||||
|
if string(result.Attachments[0].Content) != "echter anhangsinhalt" {
|
||||||
|
t.Fatalf("guter anhang hat unerwarteten inhalt: %q", result.Attachments[0].Content)
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -53,7 +53,7 @@ func (s *Session) handleSelect(ctx context.Context, cmd command) bool {
|
|||||||
}
|
}
|
||||||
|
|
||||||
mailboxName := cmd.Args[0]
|
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 {
|
if err != nil || !ok {
|
||||||
// Fehlgeschlagenes SELECT lässt den Zustand laut RFC 3501 §6.3.1
|
// Fehlgeschlagenes SELECT lässt den Zustand laut RFC 3501 §6.3.1
|
||||||
// auf Authenticated zurückfallen, nie in Selected mit ungültigem
|
// 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 {
|
if err := writeUntagged(s.writer, fmt.Sprintf("%d EXISTS", exists)); err != nil {
|
||||||
return false
|
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.state = Selected
|
||||||
s.mailbox = mailboxName
|
s.mailbox = mailboxName
|
||||||
s.mailboxSize = uint32(exists)
|
s.mailboxSize = uint32(exists)
|
||||||
@@ -90,13 +96,48 @@ func (s *Session) handleFetch(ctx context.Context, cmd command) bool {
|
|||||||
if err != nil {
|
if err != nil {
|
||||||
return s.writeErr(cmd.Tag, "NO", "FETCH failed")
|
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 {
|
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 {
|
if err := writeUntagged(s.writer, text); err != nil {
|
||||||
return false
|
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.
|
// 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
|
// tatsächliche Nachrichtenzahl des gewählten Postfachs, von SELECT
|
||||||
// gemeldet). Volle RFC-3501-Sequenzsatz-Grammatik (verschachtelte
|
// gemeldet). Volle RFC-3501-Sequenzsatz-Grammatik (verschachtelte
|
||||||
// Bereiche etc.) ist bewusst nicht Bestandteil dieser kleinsten Lösung.
|
// 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) {
|
func parseSequenceSet(raw string, maxSeq uint32) ([]uint32, error) {
|
||||||
var result []uint32
|
var result []uint32
|
||||||
for _, part := range strings.Split(raw, ",") {
|
for _, part := range strings.Split(raw, ",") {
|
||||||
|
|||||||
@@ -28,9 +28,9 @@ type fakeMailboxStore struct {
|
|||||||
mailboxes map[string][]Message
|
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]
|
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) {
|
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
|
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()) {
|
func startTestServer(t *testing.T) (addr string, stop func()) {
|
||||||
t.Helper()
|
t.Helper()
|
||||||
auth := fakeAuthenticator{users: map[string]string{"alice": "geheim123"}}
|
auth := fakeAuthenticator{users: map[string]string{"alice": "geheim123"}}
|
||||||
store := fakeMailboxStore{mailboxes: map[string][]Message{
|
store := fakeMailboxStore{mailboxes: map[string][]Message{
|
||||||
"INBOX": {
|
"INBOX": {
|
||||||
{SequenceNumber: 1, Flags: []string{"\\Seen"}},
|
{SequenceNumber: 1, UID: 101, Flags: []string{"\\Seen"}},
|
||||||
{SequenceNumber: 2, Flags: []string{}},
|
{SequenceNumber: 2, UID: 102, Flags: []string{}},
|
||||||
},
|
},
|
||||||
}}
|
}}
|
||||||
srv := NewServer(auth, store)
|
srv := NewServer(auth, store)
|
||||||
@@ -203,7 +221,7 @@ func TestCommands_AllBaseCommandsAnswered(t *testing.T) {
|
|||||||
}
|
}
|
||||||
|
|
||||||
_, lines = c.sendTagged(t, "FETCH 1 (FLAGS)")
|
_, 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)
|
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)
|
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)
|
Authenticate(ctx context.Context, username, password string) (ok bool, err error)
|
||||||
}
|
}
|
||||||
|
|
||||||
// Message ist eine minimale Nachrichtendarstellung für FETCH (nur Flags,
|
// Message ist eine minimale Nachrichtendarstellung für FETCH (nur UID +
|
||||||
// keine Inhalte — Inhaltszugriff ist Sache späterer Kacheln).
|
// 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 {
|
type Message struct {
|
||||||
SequenceNumber uint32
|
SequenceNumber uint32
|
||||||
|
UID uint32
|
||||||
Flags []string
|
Flags []string
|
||||||
}
|
}
|
||||||
|
|
||||||
// MailboxStore liefert Postfachzustand für SELECT/FETCH.
|
// MailboxStore liefert Postfachzustand für SELECT/FETCH.
|
||||||
type MailboxStore interface {
|
type MailboxStore interface {
|
||||||
// Select liefert die Anzahl der Nachrichten im Postfach mailboxName.
|
// Select liefert die Anzahl der Nachrichten sowie die UIDVALIDITY
|
||||||
// ok=false, wenn das Postfach nicht existiert.
|
// (RFC 3501 §2.3.1.1 — Pflichtbestandteil der SELECT-Antwort, Basis
|
||||||
Select(ctx context.Context, mailboxName string) (exists int, ok bool, err error)
|
// 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
|
// Fetch liefert die Nachrichten im aktuell gewählten Postfach, deren
|
||||||
// Sequenznummer in seqNumbers enthalten ist.
|
// Sequenznummer in seqNumbers enthalten ist.
|
||||||
Fetch(ctx context.Context, mailboxName string, seqNumbers []uint32) ([]Message, error)
|
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)
|
return s.handleSelect(ctx, cmd)
|
||||||
case "FETCH":
|
case "FETCH":
|
||||||
return s.handleFetch(ctx, cmd)
|
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":
|
case "LOGOUT":
|
||||||
return s.handleLogout(cmd)
|
return s.handleLogout(cmd)
|
||||||
default:
|
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,251 @@
|
|||||||
|
// IMP-04: defensive Fehlerbehandlung nicht-konformer Server. Bekannten
|
||||||
|
// Fehler vermeiden (siehe known-issues-archivmail.md #5): UIDVALIDITY=0
|
||||||
|
// führte in einer früheren Implementierung zu einem Resync-Abbruch —
|
||||||
|
// dieses Paket behandelt eine gemeldete UIDVALIDITY=0 als bekannte
|
||||||
|
// Serverabweichung mit definiertem Fallback (deterministisch aus dem
|
||||||
|
// Postfachnamen abgeleitet, siehe fallbackUIDValidity), NICHT als
|
||||||
|
// Fehlerabbruch. Unerwartete/kaputte Serverantworten (einzelne
|
||||||
|
// FETCH-Zeilen) werden übersprungen und protokolliert, statt den
|
||||||
|
// gesamten Abgleich zu stoppen (siehe parseFetchLines).
|
||||||
|
//
|
||||||
|
// Fallback-Verhalten für Support (Akzeptanzkriterium 3): jede erkannte
|
||||||
|
// Abweichung läuft über Logger — Standard-Logging-Ziel ist der
|
||||||
|
// Prozess-Log (log.Printf), bei Bedarf per WithLogger umleitbar/
|
||||||
|
// abschaltbar. Log-Präfix ist immer "imapimport: unerwartete
|
||||||
|
// server-antwort" bzw. "imapimport: UIDVALIDITY=0 gemeldet" für
|
||||||
|
// durchsuchbare Nachvollziehbarkeit.
|
||||||
|
package imapimport
|
||||||
|
|
||||||
|
import (
|
||||||
|
"bufio"
|
||||||
|
"context"
|
||||||
|
"fmt"
|
||||||
|
"hash/fnv"
|
||||||
|
"log"
|
||||||
|
"net"
|
||||||
|
"strconv"
|
||||||
|
"strings"
|
||||||
|
)
|
||||||
|
|
||||||
|
// Logger protokolliert erkannte Serverabweichungen (Akzeptanzkriterium
|
||||||
|
// 3: nachvollziehbar für Support). Signatur kompatibel mit log.Printf.
|
||||||
|
type Logger func(format string, args ...any)
|
||||||
|
|
||||||
|
func defaultLogger(format string, args ...any) {
|
||||||
|
log.Printf(format, args...)
|
||||||
|
}
|
||||||
|
|
||||||
|
// RealClient spricht echtes IMAP4rev1 (RFC 3501) über TCP — genutzt für
|
||||||
|
// den realistischen Testpostfach-Nachweis (IMP-01 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
|
||||||
|
logger Logger
|
||||||
|
}
|
||||||
|
|
||||||
|
func NewRealClient(addr, username, password string) *RealClient {
|
||||||
|
return &RealClient{addr: addr, username: username, password: password, logger: defaultLogger}
|
||||||
|
}
|
||||||
|
|
||||||
|
// WithLogger ersetzt das Standard-Logging-Ziel (z. B. für Tests, die die
|
||||||
|
// protokollierten Meldungen prüfen wollen, oder um es abzuschalten).
|
||||||
|
func (c *RealClient) WithLogger(logger Logger) *RealClient {
|
||||||
|
c.logger = logger
|
||||||
|
return c
|
||||||
|
}
|
||||||
|
|
||||||
|
func (c *RealClient) log(format string, args ...any) {
|
||||||
|
if c.logger != nil {
|
||||||
|
c.logger(format, args...)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
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 := c.resolveUIDValidity(selectLines, mailbox)
|
||||||
|
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 := c.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
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// resolveUIDValidity liest UIDVALIDITY aus der SELECT-Antwort
|
||||||
|
// (Akzeptanzkriterium 1). Eine gemeldete UIDVALIDITY=0 — bekannte
|
||||||
|
// Abweichung nicht-konformer Server (known-issues-archivmail.md #5) —
|
||||||
|
// löst einen definierten Fallback aus statt eines Abbruchs: ein
|
||||||
|
// deterministisch aus dem Postfachnamen abgeleiteter Ersatzwert, der bei
|
||||||
|
// wiederholten Läufen gegen denselben nicht-konformen Server STABIL
|
||||||
|
// bleibt (kein unnötiger Voll-Resync bei jedem einzelnen Lauf).
|
||||||
|
func (c *RealClient) resolveUIDValidity(lines []string, mailbox string) (uint64, error) {
|
||||||
|
v, found := extractUIDValidity(lines)
|
||||||
|
if !found {
|
||||||
|
c.log("imapimport: unerwartete server-antwort: keine UIDVALIDITY in SELECT-Antwort für %q gefunden, verwende fallback", mailbox)
|
||||||
|
return fallbackUIDValidity(mailbox), nil
|
||||||
|
}
|
||||||
|
if v == 0 {
|
||||||
|
c.log("imapimport: UIDVALIDITY=0 gemeldet für postfach %q (bekannte abweichung nicht-konformer server) — verwende definierten fallback statt sync-abbruch", mailbox)
|
||||||
|
return fallbackUIDValidity(mailbox), nil
|
||||||
|
}
|
||||||
|
return v, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// fallbackUIDValidity leitet einen deterministischen, garantiert von 0
|
||||||
|
// verschiedenen Ersatzwert aus dem Postfachnamen ab (FNV-1a, gleiche
|
||||||
|
// Technik wie mail/internal/search.DocumentID).
|
||||||
|
func fallbackUIDValidity(mailbox string) uint64 {
|
||||||
|
h := fnv.New64a()
|
||||||
|
_, _ = h.Write([]byte("imap-fallback-uidvalidity:"))
|
||||||
|
_, _ = h.Write([]byte(mailbox))
|
||||||
|
v := h.Sum64()
|
||||||
|
if v == 0 {
|
||||||
|
v = 1
|
||||||
|
}
|
||||||
|
return v
|
||||||
|
}
|
||||||
|
|
||||||
|
// extractUIDValidity sucht "UIDVALIDITY <n>" in den SELECT-Antwortzeilen.
|
||||||
|
// found=false, wenn keine UIDVALIDITY-Angabe vorhanden ODER sie nicht als
|
||||||
|
// Zahl lesbar ist (beides bekannte Serverabweichungen, siehe
|
||||||
|
// resolveUIDValidity — kein Fehlerabbruch an dieser Stelle).
|
||||||
|
func extractUIDValidity(lines []string) (value uint64, found bool) {
|
||||||
|
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, false
|
||||||
|
}
|
||||||
|
return v, true
|
||||||
|
}
|
||||||
|
return 0, false
|
||||||
|
}
|
||||||
|
|
||||||
|
// parseFetchLines parst Zeilen der Form
|
||||||
|
// "* <seq> FETCH (UID <uid> FLAGS (<flags>))" (siehe mail/internal/imap
|
||||||
|
// writeFetchResults). Akzeptanzkriterium 2: eine einzelne unerwartete/
|
||||||
|
// kaputte Zeile wird protokolliert und übersprungen, alle übrigen,
|
||||||
|
// korrekt lesbaren Nachrichten werden trotzdem geliefert — der gesamte
|
||||||
|
// Lauf bricht dafür NICHT ab.
|
||||||
|
func (c *RealClient) parseFetchLines(lines []string) []RemoteMessage {
|
||||||
|
var messages []RemoteMessage
|
||||||
|
for _, line := range lines {
|
||||||
|
if !strings.HasPrefix(line, "* ") {
|
||||||
|
continue // getaggte Abschlusszeile ("C3 OK ..."), kein Untagged-FETCH
|
||||||
|
}
|
||||||
|
if !strings.Contains(line, "FETCH ") {
|
||||||
|
continue // anderweitiges Untagged (z. B. künftig "* OK ..."), nichts zu parsen
|
||||||
|
}
|
||||||
|
msg, ok := parseSingleFetchLine(line)
|
||||||
|
if !ok {
|
||||||
|
c.log("imapimport: unerwartete server-antwort übersprungen: %q", line)
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
messages = append(messages, msg)
|
||||||
|
}
|
||||||
|
return messages
|
||||||
|
}
|
||||||
|
|
||||||
|
func parseSingleFetchLine(line string) (RemoteMessage, bool) {
|
||||||
|
if !strings.Contains(line, "FETCH (UID ") {
|
||||||
|
return RemoteMessage{}, false
|
||||||
|
}
|
||||||
|
uidIdx := strings.Index(line, "UID ") + len("UID ")
|
||||||
|
rest := line[uidIdx:]
|
||||||
|
spaceIdx := strings.IndexByte(rest, ' ')
|
||||||
|
if spaceIdx == -1 {
|
||||||
|
return RemoteMessage{}, false
|
||||||
|
}
|
||||||
|
uid, err := strconv.ParseUint(rest[:spaceIdx], 10, 32)
|
||||||
|
if err != nil {
|
||||||
|
return RemoteMessage{}, false
|
||||||
|
}
|
||||||
|
|
||||||
|
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, " ")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
return RemoteMessage{UID: uint32(uid), Flags: flags}, true
|
||||||
|
}
|
||||||
@@ -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,164 @@
|
|||||||
|
// IMP-04: Fehlerbehandlung nicht-konformer Server. Baut einen minimalen,
|
||||||
|
// hand-gesteuerten Fake-Server (roher TCP, KEIN mail/internal/imap) auf,
|
||||||
|
// der bewusst nicht-konforme Antworten sendet — echte Kontrolle über
|
||||||
|
// genau das Fehlerszenario, das getestet werden soll.
|
||||||
|
package imapimport
|
||||||
|
|
||||||
|
import (
|
||||||
|
"bufio"
|
||||||
|
"context"
|
||||||
|
"fmt"
|
||||||
|
"net"
|
||||||
|
"strings"
|
||||||
|
"testing"
|
||||||
|
)
|
||||||
|
|
||||||
|
// scriptedServer nimmt EINE Verbindung an und sendet exakt die
|
||||||
|
// vorgegebenen Zeilen als Antwort auf jedes eingehende Kommando (in
|
||||||
|
// Reihenfolge) — genug Kontrolle, um nicht-konforme Serverantworten
|
||||||
|
// exakt zu reproduzieren.
|
||||||
|
type scriptedServer struct {
|
||||||
|
responses [][]string // je eingehendem Kommando eine Antwortzeilen-Liste
|
||||||
|
}
|
||||||
|
|
||||||
|
func (s *scriptedServer) start(t *testing.T) (addr string) {
|
||||||
|
t.Helper()
|
||||||
|
listener, err := net.Listen("tcp", "127.0.0.1:0")
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("listener: %v", err)
|
||||||
|
}
|
||||||
|
go func() {
|
||||||
|
conn, err := listener.Accept()
|
||||||
|
if err != nil {
|
||||||
|
return
|
||||||
|
}
|
||||||
|
defer func() { _ = conn.Close() }()
|
||||||
|
reader := bufio.NewReader(conn)
|
||||||
|
|
||||||
|
_, _ = conn.Write([]byte("* OK IMAP4rev1 Service Ready\r\n"))
|
||||||
|
for _, respLines := range s.responses {
|
||||||
|
if _, err := reader.ReadString('\n'); err != nil {
|
||||||
|
return
|
||||||
|
}
|
||||||
|
for _, line := range respLines {
|
||||||
|
if _, err := conn.Write([]byte(line + "\r\n")); err != nil {
|
||||||
|
return
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}()
|
||||||
|
t.Cleanup(func() { _ = listener.Close() })
|
||||||
|
return listener.Addr().String()
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestResolveUIDValidity_ZeroTriggersDefinedFallbackNotAbort ist die
|
||||||
|
// geforderte Pflichtprüfung 1: Test simuliert Server mit UIDVALIDITY=0
|
||||||
|
// und bestätigt greifenden Fallback.
|
||||||
|
func TestResolveUIDValidity_ZeroTriggersDefinedFallbackNotAbort(t *testing.T) {
|
||||||
|
srv := &scriptedServer{responses: [][]string{
|
||||||
|
{"C1 OK LOGIN completed"},
|
||||||
|
{"* 3 EXISTS", "* OK [UIDVALIDITY 0] UIDs valid", "C2 OK [READ-WRITE] SELECT completed"},
|
||||||
|
{"* 1 FETCH (UID 1 FLAGS ())", "* 2 FETCH (UID 2 FLAGS ())", "* 3 FETCH (UID 3 FLAGS ())", "C3 OK UID FETCH completed"},
|
||||||
|
{"C4 OK LOGOUT completed"},
|
||||||
|
}}
|
||||||
|
addr := srv.start(t)
|
||||||
|
|
||||||
|
var loggedFallback bool
|
||||||
|
client := NewRealClient(addr, "user", "pass").WithLogger(func(format string, args ...any) {
|
||||||
|
msg := fmt.Sprintf(format, args...)
|
||||||
|
if strings.Contains(msg, "UIDVALIDITY=0") {
|
||||||
|
loggedFallback = true
|
||||||
|
}
|
||||||
|
})
|
||||||
|
|
||||||
|
uidvalidity, messages, err := client.Sync(context.Background(), "INBOX")
|
||||||
|
if err != nil {
|
||||||
|
// Akzeptanzkriterium 1: KEIN Sync-Abbruch bei UIDVALIDITY=0.
|
||||||
|
t.Fatalf("erwartete erfolgreichen sync trotz UIDVALIDITY=0, habe fehler: %v", err)
|
||||||
|
}
|
||||||
|
if uidvalidity == 0 {
|
||||||
|
t.Fatal("erwartete definierten fallback-wert != 0, habe weiterhin 0")
|
||||||
|
}
|
||||||
|
if len(messages) != 3 {
|
||||||
|
t.Fatalf("erwartete 3 nachrichten trotz UIDVALIDITY=0, habe %d", len(messages))
|
||||||
|
}
|
||||||
|
if !loggedFallback {
|
||||||
|
t.Fatal("erwartete protokollierten fallback-hinweis (akzeptanzkriterium 3: nachvollziehbar)")
|
||||||
|
}
|
||||||
|
|
||||||
|
// Fallback ist deterministisch für dasselbe Postfach — ein zweiter
|
||||||
|
// Aufruf gegen einen erneut nicht-konformen Server liefert real
|
||||||
|
// denselben Ersatzwert, löst also keinen unnötigen Voll-Resync bei
|
||||||
|
// jedem einzelnen Lauf aus.
|
||||||
|
if fallbackUIDValidity("INBOX") != uidvalidity {
|
||||||
|
t.Fatalf("erwartete deterministischen fallback, habe %d vs %d", fallbackUIDValidity("INBOX"), uidvalidity)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestParseFetchLines_UnexpectedResponseSkippedRestContinue ist die
|
||||||
|
// geforderte Pflichtprüfung 2: Test mit unerwarteter/kaputter
|
||||||
|
// Serverantwort bestätigt Weiterlauf für übrige Nachrichten.
|
||||||
|
func TestParseFetchLines_UnexpectedResponseSkippedRestContinue(t *testing.T) {
|
||||||
|
srv := &scriptedServer{responses: [][]string{
|
||||||
|
{"C1 OK LOGIN completed"},
|
||||||
|
{"* 3 EXISTS", "* OK [UIDVALIDITY 42] UIDs valid", "C2 OK [READ-WRITE] SELECT completed"},
|
||||||
|
{
|
||||||
|
"* 1 FETCH (UID 1 FLAGS ())",
|
||||||
|
"* GARBAGE NOT EVEN A FETCH LINE AT ALL", // kaputte/unerwartete Antwort
|
||||||
|
"* 2 FETCH SOMETHING UNPARSEABLE HERE (UID)", // ebenfalls kaputt
|
||||||
|
"* 3 FETCH (UID 3 FLAGS (\\Seen))",
|
||||||
|
"C3 OK UID FETCH completed",
|
||||||
|
},
|
||||||
|
{"C4 OK LOGOUT completed"},
|
||||||
|
}}
|
||||||
|
addr := srv.start(t)
|
||||||
|
|
||||||
|
var skippedCount int
|
||||||
|
client := NewRealClient(addr, "user", "pass").WithLogger(func(format string, args ...any) {
|
||||||
|
msg := fmt.Sprintf(format, args...)
|
||||||
|
if strings.Contains(msg, "unerwartete server-antwort übersprungen") {
|
||||||
|
skippedCount++
|
||||||
|
}
|
||||||
|
})
|
||||||
|
|
||||||
|
uidvalidity, messages, err := client.Sync(context.Background(), "INBOX")
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("erwartete erfolgreichen sync trotz kaputter zeilen, habe fehler: %v", err)
|
||||||
|
}
|
||||||
|
if uidvalidity != 42 {
|
||||||
|
t.Fatalf("erwartete uidvalidity=42, habe %d", uidvalidity)
|
||||||
|
}
|
||||||
|
// Akzeptanzkriterium 2: die BEIDEN kaputten Zeilen werden übersprungen
|
||||||
|
// UND protokolliert, die ÜBRIGEN (real 2) Nachrichten kommen trotzdem an.
|
||||||
|
if len(messages) != 2 {
|
||||||
|
t.Fatalf("erwartete 2 lesbare nachrichten trotz kaputter zeilen, habe %d: %+v", len(messages), messages)
|
||||||
|
}
|
||||||
|
if skippedCount != 2 {
|
||||||
|
t.Fatalf("erwartete 2 protokollierte übersprungene zeilen, habe %d", skippedCount)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestResolveUIDValidity_RegressionGuardAgainstZeroAbort ist die
|
||||||
|
// geforderte Pflichtprüfung 3: Regressionstest verhindert
|
||||||
|
// Wiederauftreten des UIDVALIDITY-Bugs — prüft die Fallback-Funktion
|
||||||
|
// isoliert und direkt, unabhängig vom Netzwerkpfad.
|
||||||
|
func TestResolveUIDValidity_RegressionGuardAgainstZeroAbort(t *testing.T) {
|
||||||
|
client := NewRealClient("unused:0", "u", "p")
|
||||||
|
value, err := client.resolveUIDValidity([]string{"* OK [UIDVALIDITY 0] UIDs valid"}, "INBOX")
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("regression: UIDVALIDITY=0 löste real einen fehler aus (der genau vermiedene bug): %v", err)
|
||||||
|
}
|
||||||
|
if value == 0 {
|
||||||
|
t.Fatal("regression: fallback lieferte weiterhin 0 — bug erneut aufgetreten")
|
||||||
|
}
|
||||||
|
|
||||||
|
// Fehlende UIDVALIDITY-Angabe (noch nicht-konformer als 0) darf
|
||||||
|
// ebenfalls nicht abbrechen.
|
||||||
|
value2, err := client.resolveUIDValidity([]string{"C2 OK SELECT completed"}, "INBOX")
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("regression: fehlende UIDVALIDITY löste real einen fehler aus: %v", err)
|
||||||
|
}
|
||||||
|
if value2 == 0 {
|
||||||
|
t.Fatal("regression: fallback bei fehlender UIDVALIDITY lieferte 0")
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -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
|
||||||
|
}
|
||||||
@@ -0,0 +1,132 @@
|
|||||||
|
// IMP-02: fehlertolerantes Parsing für den Import-Pfad. Additive
|
||||||
|
// Erweiterung — Parse/parseMultipart (ING-04) bleiben UNVERÄNDERT, deren
|
||||||
|
// Verhalten und Tests sind nicht Gegenstand dieser Kachel. ParseTolerant
|
||||||
|
// nutzt dieselben internen Helfer (readSinglePart, decodeTransferEncoding
|
||||||
|
// usw.), bricht aber bei EINEM fehlerhaften Teil NICHT die gesamte
|
||||||
|
// Nachricht ab (Akzeptanzkriterium 3), sondern verzeichnet den Fehler und
|
||||||
|
// verarbeitet die übrigen Teile weiter. Zusätzlich wird ein
|
||||||
|
// Gesamtgrößenlimit über alle Teile hinweg durchgesetzt
|
||||||
|
// (Akzeptanzkriterium 2 — maxAttachmentSize aus Parse ist nur das Limit
|
||||||
|
// je EINZELNEM Anhang).
|
||||||
|
package mimeparse
|
||||||
|
|
||||||
|
import (
|
||||||
|
"errors"
|
||||||
|
"fmt"
|
||||||
|
"io"
|
||||||
|
"mime"
|
||||||
|
"mime/multipart"
|
||||||
|
"net/mail"
|
||||||
|
"strings"
|
||||||
|
)
|
||||||
|
|
||||||
|
// PartError beschreibt EINEN Teil, der nicht verarbeitet werden konnte —
|
||||||
|
// die übrigen Teile der Nachricht sind davon unberührt.
|
||||||
|
type PartError struct {
|
||||||
|
Filename string
|
||||||
|
Err error
|
||||||
|
}
|
||||||
|
|
||||||
|
func (e PartError) Error() string {
|
||||||
|
return fmt.Sprintf("mimeparse: teil %q: %v", e.Filename, e.Err)
|
||||||
|
}
|
||||||
|
|
||||||
|
// ParseTolerant ist wie Parse, bricht aber bei einem fehlerhaften
|
||||||
|
// EINZELNEN Teil (z. B. überdimensionierter Anhang) nicht die gesamte
|
||||||
|
// Nachricht ab — der fehlerhafte Teil landet in den zurückgegebenen
|
||||||
|
// PartErrors, Text und übrige Anhänge werden unangetastet weiter
|
||||||
|
// verarbeitet (Akzeptanzkriterium 3). Nur eine strukturell unlesbare
|
||||||
|
// Nachricht (kaputte Kopfzeilen, fehlende Boundary) liefert weiterhin
|
||||||
|
// einen echten Fehler — davon kann sich kein Teil-für-Teil-Fallback
|
||||||
|
// erholen.
|
||||||
|
func ParseTolerant(r io.Reader, maxAttachmentSize, maxMessageSize int64) (Message, []PartError, error) {
|
||||||
|
msg, err := mail.ReadMessage(r)
|
||||||
|
if err != nil {
|
||||||
|
return Message{}, nil, fmt.Errorf("mimeparse: nachricht lesen: %w", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
mediaType, params, err := mime.ParseMediaType(msg.Header.Get("Content-Type"))
|
||||||
|
if err != nil {
|
||||||
|
body, readErr := readLimited(msg.Body, maxAttachmentSize)
|
||||||
|
if readErr != nil {
|
||||||
|
return Message{}, []PartError{{Filename: "", Err: readErr}}, nil
|
||||||
|
}
|
||||||
|
return Message{Parts: []Part{{ContentType: "text/plain", Content: body, Size: int64(len(body))}}}, nil, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
var result Message
|
||||||
|
var partErrors []PartError
|
||||||
|
budget := maxMessageSize
|
||||||
|
|
||||||
|
if strings.HasPrefix(mediaType, "multipart/") {
|
||||||
|
if err := parseMultipartTolerant(msg.Body, params["boundary"], maxAttachmentSize, &budget, &result, &partErrors); err != nil {
|
||||||
|
return Message{}, partErrors, err
|
||||||
|
}
|
||||||
|
return result, partErrors, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
part, err := readSinglePart(msg.Header.Get("Content-Transfer-Encoding"), mediaType, "", msg.Body, maxAttachmentSize)
|
||||||
|
if err != nil {
|
||||||
|
return Message{}, []PartError{{Filename: "", Err: err}}, nil
|
||||||
|
}
|
||||||
|
result.Parts = append(result.Parts, part)
|
||||||
|
return result, nil, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
func parseMultipartTolerant(r io.Reader, boundary string, maxAttachmentSize int64, budget *int64, result *Message, partErrors *[]PartError) error {
|
||||||
|
if boundary == "" {
|
||||||
|
return errors.New("mimeparse: multipart ohne boundary")
|
||||||
|
}
|
||||||
|
mr := multipart.NewReader(r, boundary)
|
||||||
|
for {
|
||||||
|
p, err := mr.NextPart()
|
||||||
|
if err == io.EOF {
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
if err != nil {
|
||||||
|
// Eine strukturell kaputte Multipart-Hülle (nicht ein
|
||||||
|
// einzelner Teil) kann von hier aus nicht sinnvoll fortgesetzt
|
||||||
|
// werden — kontrollierter Abbruch, wie in Parse.
|
||||||
|
return fmt.Errorf("mimeparse: multipart-teil lesen: %w", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
contentType := p.Header.Get("Content-Type")
|
||||||
|
mediaType, subParams, err := mime.ParseMediaType(contentType)
|
||||||
|
if err != nil {
|
||||||
|
mediaType = "text/plain"
|
||||||
|
}
|
||||||
|
filename := decodeHeaderValue(p.FileName())
|
||||||
|
|
||||||
|
if strings.HasPrefix(mediaType, "multipart/") {
|
||||||
|
if err := parseMultipartTolerant(p, subParams["boundary"], maxAttachmentSize, budget, result, partErrors); err != nil {
|
||||||
|
*partErrors = append(*partErrors, PartError{Filename: filename, Err: err})
|
||||||
|
}
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
|
||||||
|
if *budget <= 0 {
|
||||||
|
*partErrors = append(*partErrors, PartError{Filename: filename, Err: ErrMessageTooLarge})
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
|
||||||
|
part, err := readSinglePart(p.Header.Get("Content-Transfer-Encoding"), mediaType, filename, p, maxAttachmentSize)
|
||||||
|
if err != nil {
|
||||||
|
// Akzeptanzkriterium 3: NUR dieser eine Teil fällt weg,
|
||||||
|
// Verarbeitung läuft weiter.
|
||||||
|
*partErrors = append(*partErrors, PartError{Filename: filename, Err: err})
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
if part.Size > *budget {
|
||||||
|
*partErrors = append(*partErrors, PartError{Filename: filename, Err: ErrMessageTooLarge})
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
*budget -= part.Size
|
||||||
|
result.Parts = append(result.Parts, part)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// ErrMessageTooLarge wird geliefert (als PartError), wenn die Summe aller
|
||||||
|
// Anhangsgrößen einer Nachricht das Gesamtlimit überschreitet
|
||||||
|
// (Akzeptanzkriterium 2 — je-Nachricht-Limit, zusätzlich zum
|
||||||
|
// je-Anhang-Limit ErrAttachmentTooLarge aus Parse/readLimited).
|
||||||
|
var ErrMessageTooLarge = errors.New("mimeparse: nachricht überschreitet die maximal erlaubte gesamtgröße")
|
||||||
Reference in New Issue
Block a user