package backend import ( "fmt" "log" "os" "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 zuerst ins lokale mbox-Sicherheitsarchiv schreiben und diesen // Zustand festhalten; erst danach folgt TargetMailbox.Append. Beide Stufen // sind getrennt wiederaufnehmbar, damit ein Zielausfall nie das Archiv // verhindert und ein Nachlauf keine mbox-Dublette erzeugt. // 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() 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 } folderMap, err := FolderMap(a.ID) if err != nil { _ = finishJob(jobID, 0, 0, 1, "error") return err } var dst TargetMailbox var targets []TargetFolder if opened, openErr := OpenIMAPTarget(a); openErr != nil { log.Printf("%s target unavailable; continuing archive-only: %v", a.Name, openErr) } else { dst = opened targets, err = dst.Folders() if err != nil { log.Printf("%s target folders unavailable; continuing archive-only: %v", a.Name, err) _ = dst.Close() dst = nil targets = nil } } defer func() { if dst != nil { _ = dst.Close() } }() total, done, errs := 0, 0, 0 archived, targeted, targetPending := 0, 0, 0 targetUnavailableLogged := false for _, folder := range folders { if len(selectedFolders) > 0 && !selectedFolders[folder.Name] { continue } dstDelim := "/" if dst != nil { dstDelim = dst.Delim() } dstFolder := MapSourceToTarget(folder.Name, folder.Attrs, folder.Delim, dstDelim, targets, folderMap[folder.Name]) if dstFolder == "" { dstFolder = folder.Name } targetReady := false if dst != nil { if err := dst.EnsureFolder(dstFolder); err != nil { log.Printf("%s %s ensure target %s error; retrying connection: %v", a.Name, folder.Name, dstFolder, err) _ = dst.Close() dst = nil if reopened, reopenErr := OpenIMAPTarget(a); reopenErr != nil { log.Printf("%s target reconnect failed; continuing archive-only: %v", a.Name, reopenErr) } else if ensureErr := reopened.EnsureFolder(dstFolder); ensureErr != nil { log.Printf("%s %s ensure target %s after reconnect failed; continuing archive-only: %v", a.Name, folder.Name, dstFolder, ensureErr) _ = reopened.Close() } else { dst = reopened targetReady = true } } else { targetReady = true } } folderDone, folderErrs := 0, 0 folderArchived, folderTargeted, folderPending := 0, 0, 0 folderTotal := 0 identityFolder := safeMboxName(folder.Name) // After an identity-schema upgrade the indexed byte count is invalidated. // Rebuild from the actual archive before consulting copied, otherwise a // legacy Message-ID row could hide a byte-different message. archivePath := filepath.Join(mboxDir, safeMboxName(folder.Name)+mboxFileExt()) if _, statErr := os.Stat(archivePath); statErr == nil { if indexErr := EnsureMboxIndexed(archivePath); indexErr != nil { errs++ folderErrs++ log.Printf("%s %s archive identity index error: %v", a.Name, folder.Name, indexErr) continue } } else if !os.IsNotExist(statErr) { errs++ folderErrs++ log.Printf("%s %s archive stat error: %v", a.Name, folder.Name, statErr) continue } if err := src.Fetch(folder.Name, func(m RawMessage) error { total++ folderTotal++ identity := identityForRawMessage(m) state, err := GetCopyIdentityStateWithFolderAlias(a.ID, identityFolder, folder.Name, identity) if err != nil { errs++ folderErrs++ log.Printf("%s %s dedup error %s: %v", a.Name, folder.Name, m.MessageID, err) return nil } if state.MboxDone && state.TargetDone { return nil } if !state.MboxDone { 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: identityFolder, MessageID: identity.MessageID, BodySHA256: identity.BodySHA256, 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 := UpdateMboxIndexState(a.ID, identityFolder, appendInfo.FileOffset+appendInfo.FrameLen); err != nil { errs++ folderErrs++ log.Printf("%s %s mbox index state error %s: %v", a.Name, folder.Name, m.MessageID, err) return nil } if err := MarkIdentityMboxCopied(a.ID, identityFolder, identity); err != nil { errs++ folderErrs++ log.Printf("%s %s mbox stage mark error %s: %v", a.Name, folder.Name, m.MessageID, err) return nil } state.MboxDone = true archived++ folderArchived++ } if !state.TargetDone { if dst == nil || !targetReady { errs++ folderErrs++ targetPending++ folderPending++ if !targetUnavailableLogged { log.Printf("%s target unavailable; messages remain safely archived with target_done=0", a.Name) targetUnavailableLogged = true } return nil } if err := dst.Append(dstFolder, m); err != nil { log.Printf("%s %s append to %s error %s; retrying connection: %v", a.Name, folder.Name, dstFolder, m.MessageID, err) _ = dst.Close() dst = nil reopened, reopenErr := OpenIMAPTarget(a) if reopenErr == nil { reopenErr = reopened.EnsureFolder(dstFolder) } if reopenErr == nil { reopenErr = reopened.Append(dstFolder, m) } if reopenErr != nil { if reopened != nil { _ = reopened.Close() } errs++ folderErrs++ targetPending++ folderPending++ targetReady = false log.Printf("%s target reconnect/append failed; continuing archive-only: %v", a.Name, reopenErr) return nil } dst = reopened targetReady = true } if err := MarkIdentityTargetCopied(a.ID, identityFolder, identity); err != nil { errs++ folderErrs++ log.Printf("%s %s target stage mark error %s: %v", a.Name, folder.Name, m.MessageID, err) return nil } targeted++ folderTargeted++ } 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 complete=%d archived=%d target=%d pending=%d errors=%d", a.Name, folder.Name, dstFolder, folderTotal, folderDone, folderArchived, folderTargeted, folderPending, folderErrs) } state := "success" if errs > 0 { state = "partial" } _ = finishJob(jobID, total, done, errs, state) log.Printf("migration %s all folders: total=%d complete=%d archived=%d target=%d pending=%d errors=%d", a.Name, total, done, archived, targeted, targetPending, 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 }