Stream source mailbox migration

This commit is contained in:
DonVoo 2026-07-14 17:56:33 +02:00
parent f5f0ff1b8c
commit a52963c380
4 changed files with 179 additions and 33 deletions

View file

@ -1552,14 +1552,26 @@ func fetchSourceFolderMessages(account Account, folder string) ([]RawMessage, er
return nil, err return nil, err
} }
defer src.Close() defer src.Close()
return src.Fetch(folder) return collectSourceMessages(src, folder)
} }
src, err := OpenIMAPSource(account) src, err := OpenIMAPSource(account)
if err != nil { if err != nil {
return nil, err return nil, err
} }
defer src.Close() 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 { func renderInitialSourceMailboxList(account Account) string {

View file

@ -36,8 +36,8 @@ type Folder struct {
// SourceMailbox abstrahiert Quelle (IMAP primaer, POP3 Fallback), damit // SourceMailbox abstrahiert Quelle (IMAP primaer, POP3 Fallback), damit
// 07-migrate.go gegen ein Interface arbeitet. // 07-migrate.go gegen ein Interface arbeitet.
type SourceMailbox interface { type SourceMailbox interface {
Folders() ([]Folder, error) // rekursiver Ordnerbaum Folders() ([]Folder, error) // rekursiver Ordnerbaum
Fetch(folder string) ([]RawMessage, error) // BODY[] FLAGS INTERNALDATE Fetch(folder string, fn func(RawMessage) error) error // BODY[] FLAGS INTERNALDATE
Close() error Close() error
} }
@ -75,8 +75,8 @@ func (s *imapSource) Folders() ([]Folder, error) {
return folders, nil return folders, nil
} }
func (s *imapSource) Fetch(folder string) ([]RawMessage, error) { func (s *imapSource) Fetch(folder string, fn func(RawMessage) error) error {
return s.mailbox.fetchAll(folder) return s.mailbox.fetchEach(folder, fn)
} }
func (s *imapSource) Close() error { func (s *imapSource) Close() error {
@ -136,13 +136,13 @@ func (m *imapClientMailbox) listMailboxes() ([]imapListMailbox, error) {
return out, nil 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() selected, err := m.c.Select(folder, &imap.SelectOptions{ReadOnly: true}).Wait()
if err != nil { if err != nil {
return nil, err return err
} }
if selected.NumMessages == 0 { if selected.NumMessages == 0 {
return nil, nil return nil
} }
section := &imap.FetchItemBodySection{Peek: true} section := &imap.FetchItemBodySection{Peek: true}
@ -154,9 +154,7 @@ func (m *imapClientMailbox) fetchAll(folder string) ([]RawMessage, error) {
InternalDate: true, InternalDate: true,
BodySection: []*imap.FetchItemBodySection{section}, BodySection: []*imap.FetchItemBodySection{section},
}) })
defer cmd.Close()
var out []RawMessage
for { for {
data := cmd.Next() data := cmd.Next()
if data == nil { if data == nil {
@ -164,20 +162,33 @@ func (m *imapClientMailbox) fetchAll(folder string) ([]RawMessage, error) {
} }
buf, err := data.Collect() buf, err := data.Collect()
if err != nil { if err != nil {
return out, err _ = cmd.Close()
return err
} }
body := buf.FindBodySection(section) body := buf.FindBodySection(section)
if body == nil { if body == nil {
body = []byte{} body = []byte{}
} }
out = append(out, RawMessage{ if err := fn(RawMessage{
MessageID: messageIDFromFetch(buf, body), MessageID: messageIDFromFetch(buf, body),
Body: body, Body: body,
Flags: flagsToStrings(sanitizeIMAPFlags(buf.Flags)), Flags: flagsToStrings(sanitizeIMAPFlags(buf.Flags)),
InternalDate: buf.InternalDate, 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, err
} }
return out, nil return out, nil

View file

@ -123,52 +123,53 @@ func migrateAccount(a Account, selectedFolders map[string]bool) error {
if dstFolder == "" { if dstFolder == "" {
dstFolder = folder.Name 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 { if err := dst.EnsureFolder(dstFolder); err != nil {
errs++ errs++
log.Printf("%s %s ensure target %s error: %v", a.Name, folder.Name, dstFolder, err) log.Printf("%s %s ensure target %s error: %v", a.Name, folder.Name, dstFolder, err)
continue continue
} }
folderDone, folderErrs := 0, 0 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) already, err := AlreadyCopied(a.ID, folder.Name, m.MessageID)
if err != nil { if err != nil {
errs++ errs++
folderErrs++ folderErrs++
log.Printf("%s %s dedup error %s: %v", a.Name, folder.Name, m.MessageID, err) log.Printf("%s %s dedup error %s: %v", a.Name, folder.Name, m.MessageID, err)
continue return nil
} }
if already { if already {
continue return nil
} }
if err := dst.Append(dstFolder, m); err != nil { if err := dst.Append(dstFolder, m); err != nil {
errs++ errs++
folderErrs++ folderErrs++
log.Printf("%s %s append to %s error %s: %v", a.Name, folder.Name, dstFolder, m.MessageID, err) 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 { if err := mbox.Append(folder.Name, m); err != nil {
errs++ errs++
folderErrs++ folderErrs++
log.Printf("%s %s mbox error %s: %v", a.Name, folder.Name, m.MessageID, err) 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 { if err := MarkCopied(a.ID, folder.Name, m.MessageID); err != nil {
errs++ errs++
folderErrs++ folderErrs++
log.Printf("%s %s mark error %s: %v", a.Name, folder.Name, m.MessageID, err) log.Printf("%s %s mark error %s: %v", a.Name, folder.Name, m.MessageID, err)
continue return nil
} }
done++ done++
folderDone++ 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" state := "success"
if errs > 0 { if errs > 0 {
@ -209,12 +210,15 @@ func CheckAccount(name string) error {
} }
total := 0 total := 0
for _, folder := range folders { for _, folder := range folders {
msgs, err := src.Fetch(folder.Name) folderTotal := 0
if err != nil { 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) 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)) log.Printf("account %s source folder %s messages=%d", a.Name, folder.Name, folderTotal)
total += len(msgs) total += folderTotal
} }
dst, err := OpenIMAPTarget(a) dst, err := OpenIMAPTarget(a)
if err != nil { if err != nil {

119
streaming-brief.md Normal file
View file

@ -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 <konto> <src> -> <dst>: 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/<pid>/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.