diff --git a/backend/00-router.go b/backend/00-router.go index 72984b7..14f7ab4 100644 --- a/backend/00-router.go +++ b/backend/00-router.go @@ -1552,14 +1552,26 @@ func fetchSourceFolderMessages(account Account, folder string) ([]RawMessage, er return nil, err } defer src.Close() - return src.Fetch(folder) + return collectSourceMessages(src, folder) } src, err := OpenIMAPSource(account) if err != nil { return nil, err } defer src.Close() - return src.Fetch(folder) + return collectSourceMessages(src, folder) +} + +func collectSourceMessages(src SourceMailbox, folder string) ([]RawMessage, error) { + var messages []RawMessage + err := src.Fetch(folder, func(m RawMessage) error { + messages = append(messages, m) + return nil + }) + if err != nil { + return messages, err + } + return messages, nil } func renderInitialSourceMailboxList(account Account) string { diff --git a/backend/04-imap-source.go b/backend/04-imap-source.go index 04434d3..100a5d3 100644 --- a/backend/04-imap-source.go +++ b/backend/04-imap-source.go @@ -36,8 +36,8 @@ type Folder struct { // SourceMailbox abstrahiert Quelle (IMAP primaer, POP3 Fallback), damit // 07-migrate.go gegen ein Interface arbeitet. type SourceMailbox interface { - Folders() ([]Folder, error) // rekursiver Ordnerbaum - Fetch(folder string) ([]RawMessage, error) // BODY[] FLAGS INTERNALDATE + Folders() ([]Folder, error) // rekursiver Ordnerbaum + Fetch(folder string, fn func(RawMessage) error) error // BODY[] FLAGS INTERNALDATE Close() error } @@ -75,8 +75,8 @@ func (s *imapSource) Folders() ([]Folder, error) { return folders, nil } -func (s *imapSource) Fetch(folder string) ([]RawMessage, error) { - return s.mailbox.fetchAll(folder) +func (s *imapSource) Fetch(folder string, fn func(RawMessage) error) error { + return s.mailbox.fetchEach(folder, fn) } func (s *imapSource) Close() error { @@ -136,13 +136,13 @@ func (m *imapClientMailbox) listMailboxes() ([]imapListMailbox, error) { return out, nil } -func (m *imapClientMailbox) fetchAll(folder string) ([]RawMessage, error) { +func (m *imapClientMailbox) fetchEach(folder string, fn func(RawMessage) error) error { selected, err := m.c.Select(folder, &imap.SelectOptions{ReadOnly: true}).Wait() if err != nil { - return nil, err + return err } if selected.NumMessages == 0 { - return nil, nil + return nil } section := &imap.FetchItemBodySection{Peek: true} @@ -154,9 +154,7 @@ func (m *imapClientMailbox) fetchAll(folder string) ([]RawMessage, error) { InternalDate: true, BodySection: []*imap.FetchItemBodySection{section}, }) - defer cmd.Close() - var out []RawMessage for { data := cmd.Next() if data == nil { @@ -164,20 +162,33 @@ func (m *imapClientMailbox) fetchAll(folder string) ([]RawMessage, error) { } buf, err := data.Collect() if err != nil { - return out, err + _ = cmd.Close() + return err } body := buf.FindBodySection(section) if body == nil { body = []byte{} } - out = append(out, RawMessage{ + if err := fn(RawMessage{ MessageID: messageIDFromFetch(buf, body), Body: body, Flags: flagsToStrings(sanitizeIMAPFlags(buf.Flags)), InternalDate: buf.InternalDate, - }) + }); err != nil { + _ = cmd.Close() + return err + } } - if err := cmd.Close(); err != nil { + return cmd.Close() +} + +func (m *imapClientMailbox) fetchAll(folder string) ([]RawMessage, error) { + var out []RawMessage + err := m.fetchEach(folder, func(m RawMessage) error { + out = append(out, m) + return nil + }) + if err != nil { return out, err } return out, nil diff --git a/backend/07-migrate.go b/backend/07-migrate.go index e89da96..8190ac2 100644 --- a/backend/07-migrate.go +++ b/backend/07-migrate.go @@ -123,52 +123,53 @@ func migrateAccount(a Account, selectedFolders map[string]bool) error { if dstFolder == "" { dstFolder = folder.Name } - msgs, err := src.Fetch(folder.Name) - if err != nil { - errs++ - log.Printf("%s %s fetch error: %v", a.Name, folder.Name, err) - continue - } - total += len(msgs) if err := dst.EnsureFolder(dstFolder); err != nil { errs++ log.Printf("%s %s ensure target %s error: %v", a.Name, folder.Name, dstFolder, err) continue } folderDone, folderErrs := 0, 0 - for _, m := range msgs { + folderTotal := 0 + if err := src.Fetch(folder.Name, func(m RawMessage) error { + total++ + folderTotal++ already, err := AlreadyCopied(a.ID, folder.Name, m.MessageID) if err != nil { errs++ folderErrs++ log.Printf("%s %s dedup error %s: %v", a.Name, folder.Name, m.MessageID, err) - continue + return nil } if already { - continue + return nil } if err := dst.Append(dstFolder, m); err != nil { errs++ folderErrs++ log.Printf("%s %s append to %s error %s: %v", a.Name, folder.Name, dstFolder, m.MessageID, err) - continue + return nil } if err := mbox.Append(folder.Name, m); err != nil { errs++ folderErrs++ log.Printf("%s %s mbox error %s: %v", a.Name, folder.Name, m.MessageID, err) - continue + return nil } if err := MarkCopied(a.ID, folder.Name, m.MessageID); err != nil { errs++ folderErrs++ log.Printf("%s %s mark error %s: %v", a.Name, folder.Name, m.MessageID, err) - continue + return nil } done++ folderDone++ + return nil + }); err != nil { + errs++ + folderErrs++ + log.Printf("%s %s fetch error: %v", a.Name, folder.Name, err) } - log.Printf("migration %s %s -> %s: total=%d copied=%d errors=%d", a.Name, folder.Name, dstFolder, len(msgs), folderDone, folderErrs) + log.Printf("migration %s %s -> %s: total=%d copied=%d errors=%d", a.Name, folder.Name, dstFolder, folderTotal, folderDone, folderErrs) } state := "success" if errs > 0 { @@ -209,12 +210,15 @@ func CheckAccount(name string) error { } total := 0 for _, folder := range folders { - msgs, err := src.Fetch(folder.Name) - if err != nil { + folderTotal := 0 + if err := src.Fetch(folder.Name, func(m RawMessage) error { + folderTotal++ + return nil + }); err != nil { return fmt.Errorf("fetch source %s %s: %w", a.Name, folder.Name, err) } - log.Printf("account %s source folder %s messages=%d", a.Name, folder.Name, len(msgs)) - total += len(msgs) + log.Printf("account %s source folder %s messages=%d", a.Name, folder.Name, folderTotal) + total += folderTotal } dst, err := OpenIMAPTarget(a) if err != nil { diff --git a/streaming-brief.md b/streaming-brief.md new file mode 100644 index 0000000..c04a827 --- /dev/null +++ b/streaming-brief.md @@ -0,0 +1,119 @@ +# Codex-Brief — Streaming (Schritt 2 aus `go-imap-migration.md`) + +Ziel: `Fetch` puffert nicht länger **alle Bodies eines Ordners** im RAM, sondern +verarbeitet **eine Mail nach der anderen**. Das ist der **letzte Blocker vor den +echten `@dr-gold.de`-Postfächern** — deren Backup-Ordner sind noch leer, und da +reden wir über GB statt über die 199 kleinen Test-Mails. + +Aktuell: `src.Fetch(folder) ([]RawMessage, error)` lädt den kompletten Ordner in +eine Slice. Ein Postfach mit 50.000 Mails × 200 KB = **10 GB im RAM** → OOM. + +## 1. Interface-Änderung (`04-imap-source.go`) + +```go +type SourceMailbox interface { + Folders() ([]Folder, error) + Fetch(folder string, fn func(RawMessage) error) error // <— NEU: Callback + Close() error +} +``` + +Implementierung mit go-imap/v2 — **genau eine Mail gleichzeitig im Speicher**: + +```go +func (s *imapSource) Fetch(folder string, fn func(RawMessage) error) error { + sel, err := s.c.Select(folder, &imap.SelectOptions{ReadOnly: true}).Wait() + if err != nil { return err } + if sel.NumMessages == 0 { return nil } // Empty-Guard BLEIBT (kein FETCH 1:* auf leer) + + fo := &imap.FetchOptions{ + Flags: true, InternalDate: true, Envelope: true, + BodySection: []*imap.FetchItemBodySection{{Peek: true}}, + } + fcmd := s.c.Fetch(imap.SeqSetRange(1, 0), fo) // 1:* + defer fcmd.Close() + + for { + msg := fcmd.Next() + if msg == nil { break } + buf, err := msg.Collect() // puffert GENAU DIESE eine Mail + if err != nil { return err } + if err := fn(buildRawMessage(buf)); err != nil { + return err // nur FATAL -> Ordner abbrechen + } + // buf/raw gehen hier out of scope -> GC gibt den Body frei + } + return fcmd.Close() +} +``` + +## 2. `07-migrate.go` — die Pro-Mail-Schleife wandert in den Callback + +**Die Reihenfolge bleibt Zeichen für Zeichen dieselbe** — sie ist die +Idempotenz-Garantie und ist live verifiziert. Nicht umsortieren: + +```go +total, done, errs := 0, 0, 0 +err := src.Fetch(srcFolder, func(m RawMessage) error { + total++ + already, err := AlreadyCopied(a.ID, srcFolder, m.MessageID) + if err != nil { errs++; log.Printf(...); return nil } + if already { return nil } + if err := dst.Append(dstFolder, m); err != nil { errs++; log.Printf(...); return nil } + if err := mbox.Append(srcFolder, m); err != nil { errs++; log.Printf(...); return nil } + if err := MarkCopied(a.ID, srcFolder, m.MessageID); err != nil { errs++; log.Printf(...); return nil } + done++ + return nil +}) +``` + +`total` wird jetzt **während** des Streams gezählt (vorher `len(msgs)`). +Log-Zeile pro Ordner **unverändert lassen**: +`migration -> : total=N copied=M errors=E` — daran hängen +meine Prüfungen und deine Abnahme. + +## 3. Die fünf Regressions-Fallen (bitte ernst nehmen) + +1. **Der Callback darf den Ordner NICHT abbrechen.** Ein Fehler an *einer* Mail + → `errs++`, loggen, **`return nil`** (weiterstreamen). Nur ein wirklich + fataler Zustand gibt einen Fehler zurück. Sonst killt eine einzige kaputte + Mail die Migration des ganzen Ordners — heute tut sie das nicht. +2. **Im Callback NIEMALS Kommandos auf der Quell-Verbindung absetzen.** Der + FETCH läuft noch auf `s.c`; ein `Select`/`List` mittendrin zerlegt den + Protokollstrom. Der Callback fasst nur **Ziel**, **mbox** und **SQLite** an — + das ist ok, andere Verbindungen. +3. **Empty-Folder-Guard behalten** (`NumMessages == 0` → sofort `return nil`). + Der Fix darf nicht beim Refactor verloren gehen. +4. **Reihenfolge behalten:** `AlreadyCopied` → `dst.Append` → `mbox.Append` → + `MarkCopied`. Ein Crash zwischen Append und MarkCopied darf lieber doppelt + kopieren als eine Mail als „erledigt" markieren, die nie ankam. +5. **`--folders`-Scoping und `--watch` müssen weiter funktionieren** — + `RunMigrationFolders` / `selectedFolders` bleiben unverändert. + +Nebenbei: `10-pop3.go` muss die neue Interface-Signatur mitziehen (auch wenn es +noch ein Stub ist), sonst baut es nicht. + +## 4. Abnahme + +1. `go test ./...` grün. +2. **Keine Verhaltensänderung** — die bestehende Abnahme muss identisch laufen: + - `--run codex-abnahme-vdevop02-to-vdevop03` zweimal → zweiter Lauf `copied=0` + - `Gelöschte Objekte -> Papierkorb` (keine Dublette), `Entwürfe -> Entwürfe` + - Ziel-`INTERNALDATE` weiter `07-Mar-2024` (Originaldatum, nicht heute) + - Anhänge byte-identisch +3. **Der eigentliche Beweis — Speicher:** Leg in der Quelle einen Testordner mit + z. B. **200 Mails à ~1 MB** an. Alt: RSS wächst auf ~200 MB+ (ganzer Ordner + im RAM). Neu: RSS bleibt **flach** (nur eine Mail). Messen z. B. mit + `podman stats` oder `/proc//status` `VmRSS` während des Laufs. + → **Diesen Speicher-Test fahre und verifiziere ich** (ich bin Test/Review), + du musst ihn nicht selbst aufsetzen — bau nur den Code. + +## 5. NICHT in diesem Schritt + +- **Echtes Body-Streaming pro Mail** (Body als `io.Reader` direkt in APPEND + + mbox tee-en, statt `Collect()`). Damit wäre auch eine 500-MB-Einzelmail + unkritisch. Aktuell hält man **eine** Mail im RAM — das löst das gestellte + Problem (GB-Postfach) vollständig; Einzelmails sind serverseitig ohnehin + gedeckelt. Später bei Bedarf. +- POP3 verdrahten, DB-Namen-Altlast (`emailforwarder.db` vs. `mail-graveyard.db` + + zwei 0-Byte-Leichen), RFC-2047-Betreff im Forward.