Mail-Graveyard/backend/07-migrate.go
2026-07-14 20:20:30 +02:00

261 lines
7.3 KiB
Go

package backend
import (
"fmt"
"log"
"path/filepath"
"strings"
"time"
)
// RunMigration ist der Motor. Pro Konto:
// 1. Quelle oeffnen (IMAP, sonst POP3-Fallback), Ziel oeffnen. Ordnerbaum der
// Quelle (mit SPECIAL-USE + Delim) und vorhandene Ziel-Ordner holen.
// 2. Pro Quell-Ordner Zielordner via MapSourceToTarget (11-folders.go)
// bestimmen — Standardordner (Gesendet/Entwuerfe/Spam/Papierkorb/Archiv/
// Posteingang) landen im vorhandenen Rollen-Ordner des Ziels statt als
// Dublette; Pfad-Trenner wird umgehaengt. Dann EnsureFolder + Mails streamen.
// 3. Pro Mail: Message-ID gegen copied-Cache pruefen -> schon da? ueberspringen.
// Sonst DOPPELT schreiben: TargetMailbox.Append (IMAP) + MboxWriter.Append
// (lokal), dann MarkCopied.
// 4. Fortschritt in jobs-Tabelle schreiben (Browser-Anzeige).
//
// name = Konto-Name oder "all". watch = true -> Delta-Schleife bis Abbruch
// (Cutover-Fenster), sonst einmaliger Durchlauf.
func RunMigration(name string, watch bool, interval time.Duration) error {
return runMigration(name, watch, interval, nil)
}
func RunMigrationFolders(name string, folders []string) error {
return runMigration(name, false, 0, selectedFolderSet(folders))
}
func runMigration(name string, watch bool, interval time.Duration, selectedFolders map[string]bool) error {
if interval <= 0 {
interval = 60 * time.Second
}
if watch && name != "all" {
if _, err := GetAccount(name); err != nil {
return fmt.Errorf("watch account %q kann nicht geladen werden: %w", name, err)
}
}
for {
var accounts []Account
var err error
if name == "all" {
accounts, err = ListAccounts()
} else {
var a Account
a, err = GetAccount(name)
accounts = []Account{a}
}
if err != nil && (!watch || name != "all") {
return err
}
if err != nil {
log.Printf("watch load accounts failed: %v", err)
} else {
for _, a := range accounts {
if !a.Active {
continue
}
if err := migrateAccount(a, selectedFolders); err != nil {
if !watch {
return err
}
log.Printf("watch migration failed for %s: %v", a.Name, err)
}
}
}
if !watch {
return nil
}
log.Printf("watch sleeping %s", interval)
time.Sleep(interval)
}
}
// migrateAccount fuehrt einen einzelnen Umzug aus (ein Durchlauf).
func migrateAccount(a Account, selectedFolders map[string]bool) error {
if a.ID == 0 {
stored, err := GetAccount(a.Name)
if err != nil {
return err
}
a = stored
}
src, err := OpenIMAPSource(a)
if err != nil {
return fmt.Errorf("open source %s: %w", a.Name, err)
}
defer src.Close()
dst, err := OpenIMAPTarget(a)
if err != nil {
return fmt.Errorf("open target %s: %w", a.Name, err)
}
defer dst.Close()
archiveName := accountMboxDir(a)
if archiveName == "" {
return fmt.Errorf("keine Archiv-Mailbox fuer %s ausgewaehlt", a.Name)
}
mboxDir := filepath.Join(Cfg.MboxRoot, archiveName)
mbox, err := NewMboxWriter(mboxDir)
if err != nil {
return err
}
jobID, _ := startJob(a.ID)
folders, err := src.Folders()
if err != nil {
_ = finishJob(jobID, 0, 0, 1, "error")
return err
}
targets, err := dst.Folders()
if err != nil {
_ = finishJob(jobID, 0, 0, 1, "error")
return err
}
total, done, errs := 0, 0, 0
for _, folder := range folders {
if len(selectedFolders) > 0 && !selectedFolders[folder.Name] {
continue
}
dstFolder := MapSourceToTarget(folder.Name, folder.Attrs, folder.Delim, dst.Delim(), targets, "")
if dstFolder == "" {
dstFolder = folder.Name
}
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
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)
return nil
}
if already {
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)
return nil
}
appendInfo, err := mbox.Append(folder.Name, m)
if err != nil {
errs++
folderErrs++
log.Printf("%s %s mbox error %s: %v", a.Name, folder.Name, m.MessageID, err)
return nil
}
if err := SaveMboxIndex(MboxIndexEntry{
AccountID: a.ID,
Folder: folder.Name,
MessageID: m.MessageID,
Subject: appendInfo.Subject,
From: appendInfo.From,
Date: appendInfo.Date,
FileOffset: appendInfo.FileOffset,
FrameLen: appendInfo.FrameLen,
InnerOffset: appendInfo.InnerOffset,
InnerLen: appendInfo.InnerLen,
}); err != nil {
errs++
folderErrs++
log.Printf("%s %s mbox index error %s: %v", a.Name, folder.Name, m.MessageID, err)
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)
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, folderTotal, folderDone, folderErrs)
}
state := "success"
if errs > 0 {
state = "partial"
}
_ = finishJob(jobID, total, done, errs, state)
log.Printf("migration %s all folders: total=%d copied=%d errors=%d", a.Name, total, done, errs)
if errs > 0 {
return fmt.Errorf("migration %s finished with %d errors", a.Name, errs)
}
return nil
}
func selectedFolderSet(folders []string) map[string]bool {
out := make(map[string]bool, len(folders))
for _, folder := range folders {
folder = strings.TrimSpace(folder)
if folder != "" {
out[folder] = true
}
}
return out
}
func CheckAccount(name string) error {
a, err := GetAccount(name)
if err != nil {
return err
}
src, err := OpenIMAPSource(a)
if err != nil {
return fmt.Errorf("open source %s: %w", a.Name, err)
}
defer src.Close()
folders, err := src.Folders()
if err != nil {
return fmt.Errorf("list source folders %s: %w", a.Name, err)
}
total := 0
for _, folder := range folders {
folderTotal, err := src.Count(folder.Name)
if err != nil {
return fmt.Errorf("count source %s %s: %w", a.Name, folder.Name, err)
}
log.Printf("account %s source folder %s messages=%d", a.Name, folder.Name, folderTotal)
total += folderTotal
}
dst, err := OpenIMAPTarget(a)
if err != nil {
return fmt.Errorf("open target %s: %w", a.Name, err)
}
defer dst.Close()
log.Printf("account %s ok: source folders=%d messages=%d, target login ok", a.Name, len(folders), total)
return nil
}
func startJob(accountID int64) (int64, error) {
res, err := DB.Exec(`INSERT INTO jobs(account_id,state) VALUES(?, 'running')`, accountID)
if err != nil {
return 0, err
}
return res.LastInsertId()
}
func finishJob(jobID int64, total, done, errs int, state string) error {
if jobID == 0 {
return nil
}
_, err := DB.Exec(`UPDATE jobs SET finished=CURRENT_TIMESTAMP,total=?,done=?,errors=?,state=? WHERE id=?`, total, done, errs, state, jobID)
return err
}