Add CLI INBOX migration breakthrough
This commit is contained in:
parent
8a6160817b
commit
27206930f0
6 changed files with 713 additions and 43 deletions
|
|
@ -1,5 +1,12 @@
|
|||
package backend
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"log"
|
||||
"path/filepath"
|
||||
"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.
|
||||
|
|
@ -15,14 +22,120 @@ package backend
|
|||
// name = Konto-Name oder "all". watch = true -> Delta-Schleife bis Abbruch
|
||||
// (Cutover-Fenster), sonst einmaliger Durchlauf.
|
||||
func RunMigration(name string, watch bool) error {
|
||||
// TODO Codex: Accounts laden (name/"all"), pro Konto migrateAccount().
|
||||
// Bei watch: Schleife mit Pause, bis Kontext abgebrochen; dank copied-Cache
|
||||
// kopiert jeder Durchlauf nur Neues.
|
||||
return nil
|
||||
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 {
|
||||
return err
|
||||
}
|
||||
for _, a := range accounts {
|
||||
if !a.Active {
|
||||
continue
|
||||
}
|
||||
if err := migrateAccount(a); err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
if !watch {
|
||||
return nil
|
||||
}
|
||||
time.Sleep(60 * time.Second)
|
||||
}
|
||||
}
|
||||
|
||||
// migrateAccount fuehrt einen einzelnen Umzug aus (ein Durchlauf).
|
||||
func migrateAccount(a Account) error {
|
||||
// TODO Codex: siehe Schritte oben. Fehler pro Mail zaehlen, nicht abbrechen.
|
||||
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()
|
||||
mboxDir := filepath.Join(Cfg.MboxRoot, a.MboxDir)
|
||||
mbox, err := NewMboxWriter(mboxDir)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
jobID, _ := startJob(a.ID)
|
||||
msgs, err := src.Fetch("INBOX")
|
||||
if err != nil {
|
||||
_ = finishJob(jobID, 0, 0, 1, "error")
|
||||
return err
|
||||
}
|
||||
if err := dst.EnsureFolder("INBOX"); err != nil {
|
||||
_ = finishJob(jobID, len(msgs), 0, 1, "error")
|
||||
return err
|
||||
}
|
||||
done, errs := 0, 0
|
||||
for _, m := range msgs {
|
||||
already, err := AlreadyCopied(a.ID, "INBOX", m.MessageID)
|
||||
if err != nil {
|
||||
errs++
|
||||
log.Printf("%s INBOX dedup error %s: %v", a.Name, m.MessageID, err)
|
||||
continue
|
||||
}
|
||||
if already {
|
||||
continue
|
||||
}
|
||||
if err := dst.Append("INBOX", m); err != nil {
|
||||
errs++
|
||||
log.Printf("%s INBOX append error %s: %v", a.Name, m.MessageID, err)
|
||||
continue
|
||||
}
|
||||
if err := mbox.Append("INBOX", m); err != nil {
|
||||
errs++
|
||||
log.Printf("%s INBOX mbox error %s: %v", a.Name, m.MessageID, err)
|
||||
continue
|
||||
}
|
||||
if err := MarkCopied(a.ID, "INBOX", m.MessageID); err != nil {
|
||||
errs++
|
||||
log.Printf("%s INBOX mark error %s: %v", a.Name, m.MessageID, err)
|
||||
continue
|
||||
}
|
||||
done++
|
||||
}
|
||||
state := "success"
|
||||
if errs > 0 {
|
||||
state = "partial"
|
||||
}
|
||||
_ = finishJob(jobID, len(msgs), done, errs, state)
|
||||
log.Printf("migration %s INBOX: total=%d copied=%d errors=%d", a.Name, len(msgs), done, errs)
|
||||
if errs > 0 {
|
||||
return fmt.Errorf("migration %s finished with %d errors", a.Name, errs)
|
||||
}
|
||||
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
|
||||
}
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue